Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bf1937bc4f | ||
|
|
6dffa3daff | ||
|
|
8891b27cec | ||
|
|
f373caae67 | ||
|
|
f27e0eed50 | ||
|
|
0e488f27fd | ||
|
|
06173cea47 | ||
|
|
333892725a | ||
|
|
bb8f0ec1ba | ||
|
|
4e1304df3a | ||
|
|
3d6d487b59 | ||
|
|
3cc962c3d0 | ||
|
|
a375d09884 | ||
|
|
ac9b55d7e4 | ||
|
|
92dea5889b | ||
|
|
e4fec46afc | ||
|
|
6454dca048 | ||
|
|
eed02aaddc | ||
|
|
e09f96466c | ||
|
|
52381157d9 | ||
|
|
bb83b0db05 | ||
|
|
0d4f59bf68 | ||
|
|
c1b8aeb692 | ||
|
|
6277f6d83f | ||
|
|
1247498e57 | ||
|
|
693c03a16c | ||
|
|
62029e4a30 | ||
|
|
ca97c40d9f | ||
|
|
230e25dbf0 | ||
|
|
9347df2159 | ||
|
|
0d582b6458 | ||
|
|
2878a6fcdf | ||
|
|
367ea64dcd | ||
|
|
a6c05d3b48 | ||
|
|
4bfd08cead | ||
|
|
c75539c167 | ||
|
|
962aa5b573 | ||
|
|
068e055501 | ||
|
|
feb9c12bd2 | ||
|
|
5b10d3dcfb | ||
|
|
a58907ed7f | ||
|
|
8a12a9b5fa | ||
|
|
38cad77c1e | ||
|
|
9daee2c6cc | ||
|
|
4a30cc705b | ||
|
|
93ca8b49c3 | ||
|
|
c30ca2dba4 | ||
|
|
e7aed0e765 | ||
|
|
400fc272eb | ||
|
|
91ba706b55 | ||
|
|
abbd1f5aaa | ||
|
|
16baab39ef | ||
|
|
a1624ca69c | ||
|
|
41bae638fc | ||
|
|
b7166b6283 | ||
|
|
6abca7db24 | ||
|
|
e935bad9ac | ||
|
|
9617b568ed | ||
|
|
146b79cf74 | ||
|
|
174863596f | ||
|
|
45cb9e3ce6 | ||
|
|
938be883ad | ||
|
|
5abbce1512 | ||
|
|
3264df7836 | ||
|
|
62e368f812 | ||
|
|
f5242b0488 | ||
|
|
9e50767bf9 | ||
|
|
b1a5a3739e | ||
|
|
8c770b0313 | ||
|
|
27175ceb09 | ||
|
|
bfdc3ef4dd | ||
|
|
2f9d3df5aa | ||
|
|
01eb014aeb | ||
|
|
06aa2c6d68 | ||
|
|
fc5dd9e2df | ||
|
|
3ac04864e8 | ||
|
|
0e9f7dc860 | ||
|
|
c6459d8448 | ||
|
|
54198b9679 | ||
|
|
d5b527a0eb | ||
|
|
8677e0446f | ||
|
|
5f72b80651 | ||
|
|
c41cf14996 | ||
|
|
b335f98143 | ||
|
|
09803955b2 | ||
|
|
61026459c9 | ||
|
|
0079d8700f | ||
|
|
803aa1d88d | ||
|
|
dc343856ba | ||
|
|
94f24ce096 | ||
|
|
f6c3eb6e89 | ||
|
|
95f3dd41e1 | ||
|
|
2356537a5d | ||
|
|
e0ceef7b8d | ||
|
|
fe9132bd76 | ||
|
|
5bccfb9adc | ||
|
|
052c9d5eb6 | ||
|
|
30b3e11325 | ||
|
|
a0f5bbf1d4 | ||
|
|
181781247f | ||
|
|
9e9b4b3392 | ||
|
|
0f97cfc01e | ||
|
|
dca5255f70 | ||
|
|
e3555bab49 | ||
|
|
e757c308b6 | ||
|
|
d325fe2a98 | ||
|
|
73e663c18a | ||
|
|
2ab442fd2e | ||
|
|
e76578beb9 | ||
|
|
74e5c65afc | ||
|
|
06bd72103b | ||
|
|
d16fc5bc96 | ||
|
|
44ace4c9f5 | ||
|
|
f90d348217 | ||
|
|
79bc15cf7c | ||
|
|
d4718c9e50 | ||
|
|
50d4675884 | ||
|
|
fa850833a5 |
@@ -22,7 +22,7 @@ jobs:
|
||||
id-token: write
|
||||
environment:
|
||||
name: pypi
|
||||
url: https://pypi.org/project/comfyui-idl/${{ github.ref_name }}
|
||||
url: https://pypi.org/project/comfy-pack/${{ github.ref_name }}
|
||||
runs-on: ubuntu-latest
|
||||
needs: build-package
|
||||
|
||||
|
||||
@@ -1,73 +1,207 @@
|
||||
# ComfyPack
|
||||
# Comfy-Pack: Making ComfyUI Workflows Shareable
|
||||
|
||||
A comprehensive toolkit for standardizing, packaging and deploying ComfyUI workflows as reproducible environments and production-ready REST services.
|
||||

|
||||
|
||||
## Features
|
||||
- **Package Everything**: Create reproducible `.cpack.zip` files containing your workflow, custom nodes, model versions, and all dependencies
|
||||
- **Standardize Parameters**: Define and validate workflow inputs through UI nodes for images, text, numbers and more
|
||||
- **CLI Support**: Restore environment and run inference from command line
|
||||
- **REST API Generation**: Auto-convert any workflow into REST service with OpenAPI docs
|
||||
|
||||
## Quick Start
|
||||
comfy-pack is a comprehensive toolkit for reliably packing and unpacking environments for ComfyUI workflows.
|
||||
|
||||
|
||||
- 📦 **Pack workflow environments as artifacts:** Saves the workflow environment in a `.cpack.zip` artifact with Python package versions, ComfyUI and custom node revisions, and model hashes.
|
||||
- ✨ **Unpack artifacts to recreate workflow environments:** Unpacks the `.cpack.zip` artifact to recreate the same environment with the exact Python package versions, ComfyUI and custom node revisions, and model weights.
|
||||
- 🚀 **Deploy workflows as APIs:** Deploys the workflow as a RESTful API with customizable input and output parameters.
|
||||
|
||||
## Motivations
|
||||
ComfyUI Manager is great for find missing custom nodes. But when sharing ComfyUI workflows to others(your audience or team members), you've still likely heard these responses:
|
||||
|
||||
- "Custom Node not found"
|
||||
- "Cannot find the correct model file"
|
||||
- "Missing Python dependencies"
|
||||
|
||||
These are fundamental challenges in workflow sharing – every component should match exactly: custom nodes, model files, and Python dependencies. Modern pacakge managers like npm and poetry introduced "lock" feature, which means record the exact version for every requirement. ComfyUI Manager isn't designed for that.
|
||||
|
||||
We learned it from our community and developed comfy-pack to address these problems. With a single click, it captures and locks your entire workflow environment into a `.cpack.zip` file, including Python packages, custom nodes, model hashes, and required assets.
|
||||
|
||||
Users can recreate the exact environment with one command:
|
||||
|
||||
```bash
|
||||
comfy-pack unpack workflow.cpack.zip
|
||||
```
|
||||
|
||||
This means you can focus on your creative work while comfy-pack handles the rest.
|
||||
|
||||
## Usages
|
||||
|
||||
### Installation
|
||||
|
||||
We recommend you use ComfyUI Manager to install comfy-pack. Simply search for `comfy-pack` and click **Install**. Restart the server and refresh your ComfyUI interface to apply changes.
|
||||
|
||||

|
||||
|
||||
Alternatively, clone the project repository through `git`.
|
||||
|
||||
```bash
|
||||
cd ComfyUI/custom_nodes
|
||||
git clone https://github.com/bentoml/comfy-pack.git
|
||||
```
|
||||
|
||||
To install the comfy-pack CLI, run:
|
||||
|
||||
```bash
|
||||
pip install comfy-pack
|
||||
```
|
||||
|
||||
### Create a Pack
|
||||
1. Install ComfyPack custom nodes in ComfyUI
|
||||
2. Design your workflow with parameter nodes
|
||||
3. Click "Package" button to create `.cpack.zip`
|
||||
### Pack a ComfyUI workflow and its environment
|
||||
|
||||
You can package a workflow and the environment required to run the workflow into an artifact that can be unpacked elsewhere.
|
||||
|
||||
1. Click the **Package** button to create a `.cpack.zip` artifact.
|
||||
2. (Optional) Select the models that you want to include (only model hash will be recorded, so you won't get a 100GB zip file).
|
||||
|
||||

|
||||
|
||||
### Unpack the ComfyUI environments
|
||||
|
||||
Unpacking a `.cpack.zip` artifact will restore the ComfyUI environment for the workflow. During unpacking, comfy-pack will perform the following steps.
|
||||
|
||||
1. Prepare a Python virtual environment with the exact packages used to run the workflow.
|
||||
2. Clone ComfyUI and custom nodes from the exact revisions required by the workflow.
|
||||
3. Search for and download models from common registries like Hugging Face and Civitai. Unpacking workflows using the same model will not cause the model to be downloaded multiple times. Instead, model weights will be symbolically linked.
|
||||
|
||||
To unpack:
|
||||
|
||||
### Restore to a ComfyUI project
|
||||
```bash
|
||||
# Restore environment from pack, will install everything needed except the models.
|
||||
comfy-pack restore workflow.cpack.zip --dir ./
|
||||
comfy-pack unpack workflow.cpack.zip
|
||||
```
|
||||
|
||||
### Run Inference
|
||||
For example cpack files, check our [examples folder](examples/).
|
||||
|
||||
### Deploy a workflow as an API
|
||||
|
||||
You can turn a ComfyUI workflow into an API endpoint callable using any clients through HTTP.
|
||||
|
||||
<details>
|
||||
<summary> 1. Annotate input & output </summary>
|
||||
|
||||
Use custom nodes provided by comfy-pack to annotate the fields to be used as input and output parameters. To add a comfy-pack node, right-click and select **Add Node** > **ComfyPack** > **output/input** > [Select a type]
|
||||
|
||||
Input nodes:
|
||||
|
||||
- ImageInput: Accepts `image` type input, similar to the official `LoadImage` node
|
||||
- StringInput: Accepts `string` type input (e.g., prompts)
|
||||
- IntInput: Accepts `int` type input (e.g., dimensions, seeds)
|
||||
- AnyInput: Accepts `combo` type and more input (e.g., custom nodes)
|
||||
|
||||

|
||||
|
||||
Output nodes:
|
||||
|
||||
- ImageOutput: Outputs `image` type, similar to the official `SaveImage` node
|
||||
- FileOutput: Outputs file path as `string` type and saves the file under that path
|
||||
|
||||

|
||||
|
||||
More field types are under way.
|
||||
</details>
|
||||
|
||||
<details>
|
||||
<summary> 2. Serve the workflow </summary>
|
||||
|
||||
Start an HTTP server at `http://127.0.0.1:3000` (default) to serve the workflow under the `/generate` path.
|
||||
|
||||

|
||||
|
||||
You can call the `/generate` endpoint by specifying parameters configured through your comfy-pack nodes, such as prompt, width, height, and seed.
|
||||
|
||||
> [!NOTE]
|
||||
> The name of a comfy-pack node is the parameter name used for API calls.
|
||||
|
||||
Examples to call the endpoint:
|
||||
|
||||
CURL
|
||||
|
||||
```bash
|
||||
curl -X 'POST' \
|
||||
'http://127.0.0.1:3000/generate' \
|
||||
-H 'accept: application/octet-stream' \
|
||||
-H 'Content-Type: application/json' \
|
||||
-d '{
|
||||
"prompt": "rocks in a bottle",
|
||||
"width": 512,
|
||||
"height": 512,
|
||||
"seed": 1
|
||||
}'
|
||||
```
|
||||
|
||||
BentoML client
|
||||
|
||||
Under the hood, comfy-pack leverages [BentoML](https://github.com/bentoml/BentoML), the unified model serving framework. You can invoke the endpoint using [the BentoML Python client](https://docs.bentoml.com/en/latest/build-with-bentoml/clients.html):
|
||||
|
||||
```python
|
||||
import bentoml
|
||||
|
||||
with bentoml.SyncHTTPClient("http://127.0.0.1:3000") as client:
|
||||
result = client.generate(
|
||||
prompt="rocks in a bottle",
|
||||
width=512,
|
||||
height=512,
|
||||
seed=1
|
||||
)
|
||||
```
|
||||
|
||||
</details>
|
||||
|
||||
<details>
|
||||
<summary> 3. (Optional) Pack the workflow and environment </summary>
|
||||
|
||||
Pack the workflow and environment into an artifact that can be unpacked elsewhere to recreate the workflow.
|
||||
|
||||
```bash
|
||||
# Get the workflow input spec
|
||||
comfy-pack info workflow.cpack.zip
|
||||
comfy-pack run workflow.cpack.zip --help
|
||||
|
||||
# Run
|
||||
comfy-pack run workflow.cpack.zip --src-image image.png --video video.mp4
|
||||
comfy-pack run workflow.cpack.zip --src-image image.png --video video.mp4
|
||||
```
|
||||
</details>
|
||||
|
||||
### Online REST service
|
||||
under development
|
||||
<details>
|
||||
<summary> 4. (Optional) Deploy to the cloud </summary>
|
||||
|
||||
## Parameter Nodes
|
||||
Deploy to [BentoCloud](https://www.bentoml.com/) with access to a variety of GPUs and blazing fast scaling.
|
||||
|
||||
ComfyPack provides custom nodes for standardizing inputs:
|
||||
- ImageInput
|
||||
- StringInput
|
||||
- IntInput
|
||||
- AnyInput
|
||||
- ImageOutput
|
||||
- FileOutput
|
||||
- ...
|
||||
Follow [the instructions here](https://docs.bentoml.com/en/latest/scale-with-bentocloud/manage-api-tokens.html) to get your BentoCloud access token. If you don’t have a BentoCloud account, you can [sign up for free](https://bentoml.com/).
|
||||
|
||||
These nodes help define clear interfaces for your workflow.
|
||||

|
||||
|
||||
## Docker Support
|
||||
Under development
|
||||
</details>
|
||||
|
||||
## Security Guidelines
|
||||
|
||||
A cpack file only contains the metadata of the workflow environment, such as Python package versions, ComfyUI and custom node revisions, and model hashes. It does not contain any sensitive information like API keys, passwords, or user data. However, unpacking a cpack file will install custom nodes and Python dependencies. It is recommended to unpack cpack files from trusted sources.
|
||||
|
||||
comfy-pack has a strict mode for unpacking. You can enable it by setting the `CPACK_STRICT_MODE` environment variable to `true`. It will sacrifice some flexibility and compatibility for security. For now, comfy-pack will:
|
||||
|
||||
* Use more strict index strategy in Python package installation
|
||||
|
||||
More security features are under way.
|
||||
|
||||
|
||||
## Examples
|
||||
## Roadmap
|
||||
|
||||
Check our [examples folder](examples/) for:
|
||||
- Basic workflow packaging
|
||||
- Parameter configuration
|
||||
- API integration
|
||||
- Docker deployment
|
||||
This project is under active development. Currently we are working on:
|
||||
|
||||
## License
|
||||
MIT License
|
||||
- Enhanced user experience
|
||||
- Docker support
|
||||
- Local `.cpack` file management with version control
|
||||
- Enhanced service capabilities
|
||||
|
||||
## Community
|
||||
- Issues & Feature Requests: GitHub Issues
|
||||
- Questions & Discussion: Discord Server
|
||||
|
||||
Detailed documentation: under development
|
||||
comfy-pack is actively maintained by the BentoML team. Feel free to reach out 👉 [Join our Slack community!](https://l.bentoml.com/join-slack)
|
||||
|
||||
## Contributing
|
||||
|
||||
As an open-source project, we welcome contributions of all kinds, such as new features, bug fixes, and documentation. Here are some of the ways to contribute:
|
||||
|
||||
- Repost a bug by creating a [GitHub issue](https://github.com/bentoml/comfy-pack/issues).
|
||||
- Submit a [pull request](https://github.com/bentoml/comfy-pack/pulls) or help review other developers’ pull requests.
|
||||
|
||||
Binary file not shown.
@@ -0,0 +1,8 @@
|
||||
import sys
|
||||
import pathlib
|
||||
|
||||
|
||||
SRC_DIR = pathlib.Path(__file__).parent.parent / "src"
|
||||
|
||||
if str(SRC_DIR) not in sys.path:
|
||||
sys.path.append(str(SRC_DIR))
|
||||
|
||||
+356
-108
@@ -1,43 +1,58 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
import shutil
|
||||
import socket
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
import time
|
||||
import uuid
|
||||
import zipfile
|
||||
from importlib.metadata import Distribution, distributions
|
||||
from pathlib import Path
|
||||
from typing import Union
|
||||
from typing import Any, Union
|
||||
|
||||
import folder_paths
|
||||
from aiohttp import web
|
||||
from server import PromptServer
|
||||
|
||||
from comfy_pack.hash import async_batch_get_sha256
|
||||
from comfy_pack.model_helper import alookup_model_source
|
||||
from comfy_pack.package import build_bento
|
||||
|
||||
ZPath = Union[Path, zipfile.Path]
|
||||
TEMP_FOLDER = Path(__file__).parent.parent / "temp"
|
||||
COMFY_PACK_DIR = Path(__file__).parent.parent / "src/comfy_pack"
|
||||
EXCLUDE_PACKAGES = ["bentoml", "onnxruntime"] # TODO: standardize this
|
||||
COMFY_PACK_DIR = Path(__file__).parent.parent / "src" / "comfy_pack"
|
||||
EXCLUDE_PACKAGES = ["bentoml", "onnxruntime", "conda"] # TODO: standardize this
|
||||
|
||||
|
||||
def _get_requirement_string(dist: Distribution) -> str:
|
||||
direct_url_text = dist.read_text("direct_url.json")
|
||||
pinned_str = f'{dist.metadata["Name"]}=={dist.version}'
|
||||
if not direct_url_text:
|
||||
return pinned_str
|
||||
direct_url = json.loads(direct_url_text)
|
||||
if url := direct_url.get("url"):
|
||||
if url.startswith("file://"):
|
||||
# we are not able to share local files
|
||||
return pinned_str
|
||||
if vcs_info := direct_url.get("vcs_info"):
|
||||
url = f"{vcs_info['vcs']}+{url}@{vcs_info['commit_id']}"
|
||||
if subdirectory := direct_url.get("subdirectory"):
|
||||
url += f"#subdirectory={subdirectory}"
|
||||
return f"{dist.metadata['Name']} @ {url}"
|
||||
else:
|
||||
return pinned_str
|
||||
|
||||
|
||||
async def _write_requirements(path: ZPath, extras: list[str] | None = None) -> None:
|
||||
print("Package => Writing requirements.txt")
|
||||
with path.joinpath("requirements.txt").open("w") as f:
|
||||
proc = await asyncio.subprocess.create_subprocess_exec(
|
||||
sys.executable,
|
||||
"-m",
|
||||
"pip",
|
||||
"freeze",
|
||||
"--exclude-editable",
|
||||
*[f"--exclude={p}" for p in EXCLUDE_PACKAGES],
|
||||
stdout=subprocess.PIPE,
|
||||
)
|
||||
stdout, _ = await proc.communicate()
|
||||
f.write(stdout.decode().rstrip("\n") + "\n")
|
||||
for dist in distributions():
|
||||
f.write(_get_requirement_string(dist) + "\n")
|
||||
if extras:
|
||||
f.write("\n".join(extras) + "\n")
|
||||
|
||||
@@ -48,7 +63,8 @@ async def _write_snapshot(path: ZPath, data: dict, models: list | None = None) -
|
||||
)
|
||||
stdout, _ = await proc.communicate()
|
||||
if models is None:
|
||||
models = await _get_models(data)
|
||||
print("Package => Writing models")
|
||||
models = await _get_models()
|
||||
with path.joinpath("snapshot.json").open("w") as f:
|
||||
data = {
|
||||
"python": f"{sys.version_info.major}.{sys.version_info.minor}",
|
||||
@@ -59,44 +75,42 @@ async def _write_snapshot(path: ZPath, data: dict, models: list | None = None) -
|
||||
f.write(json.dumps(data, indent=2))
|
||||
|
||||
|
||||
def _get_file_info_key(file_path: Path) -> tuple:
|
||||
# use mtime, size to determine if a file has changed
|
||||
return (file_path.stat().st_mtime, file_path.stat().st_size)
|
||||
def _is_port_in_use(port: int | str, host="localhost"):
|
||||
if isinstance(port, str):
|
||||
port = int(port)
|
||||
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
|
||||
try:
|
||||
s.connect((host, port))
|
||||
return True
|
||||
except ConnectionRefusedError:
|
||||
return False
|
||||
except Exception:
|
||||
return True
|
||||
|
||||
|
||||
def _get_hash(file_path: Path) -> str:
|
||||
cpack_hash_cache_file = TEMP_FOLDER / "model_sha_cache.json"
|
||||
try:
|
||||
with cpack_hash_cache_file.open("r") as f:
|
||||
cache = json.load(f)
|
||||
except (FileNotFoundError, json.JSONDecodeError):
|
||||
cache = {}
|
||||
|
||||
mtime, size = _get_file_info_key(file_path)
|
||||
if file_path.absolute().as_posix() in cache:
|
||||
if (
|
||||
cache[file_path.absolute().as_posix()]["mtime"] == mtime
|
||||
and cache[file_path.absolute().as_posix()]["size"] == size
|
||||
):
|
||||
return cache[file_path.absolute().as_posix()]["sha256"]
|
||||
|
||||
with cpack_hash_cache_file.open("w") as f:
|
||||
hasher = hashlib.sha256()
|
||||
with open(file_path, "rb") as file:
|
||||
for chunk in iter(lambda: file.read(4096), b""):
|
||||
hasher.update(chunk)
|
||||
sha = hasher.hexdigest()
|
||||
cache[file_path.absolute().as_posix()] = {
|
||||
"sha256": sha,
|
||||
"mtime": mtime,
|
||||
"size": size,
|
||||
}
|
||||
json.dump(cache, f)
|
||||
return sha
|
||||
def _is_file_refered(file_path: Path, workflow_api: dict) -> bool:
|
||||
""" """
|
||||
used_inputs = set()
|
||||
for node in workflow_api.values():
|
||||
for _, v in node["inputs"].items():
|
||||
if isinstance(v, str):
|
||||
used_inputs.add(v)
|
||||
all_inputs = "\n".join(used_inputs)
|
||||
file_path = file_path.absolute().relative_to(folder_paths.base_path)
|
||||
if file_path.parts[0] == "input":
|
||||
relpath = Path(*file_path.parts[1:])
|
||||
else: # models
|
||||
relpath = Path(*file_path.parts[2:])
|
||||
return str(relpath) in all_inputs
|
||||
|
||||
|
||||
async def _get_models(data: dict, store_models: bool = False) -> list:
|
||||
print("Package => Writing models")
|
||||
async def _get_models(
|
||||
store_models: bool = False,
|
||||
workflow_api: dict | None = None,
|
||||
model_filter: set[str] | None = None,
|
||||
ensure_sha=True,
|
||||
ensure_source=True,
|
||||
) -> list:
|
||||
proc = await asyncio.subprocess.create_subprocess_exec(
|
||||
"git",
|
||||
"ls-files",
|
||||
@@ -106,28 +120,42 @@ async def _get_models(data: dict, store_models: bool = False) -> list:
|
||||
)
|
||||
stdout, _ = await proc.communicate()
|
||||
|
||||
used_inputs = set()
|
||||
for node in data["workflow_api"].values():
|
||||
for _, v in node["inputs"].items():
|
||||
if isinstance(v, str):
|
||||
used_inputs.add(v)
|
||||
|
||||
models = []
|
||||
for line in stdout.decode().splitlines():
|
||||
if os.path.basename(line).startswith("."):
|
||||
continue
|
||||
filename = os.path.abspath(line)
|
||||
relpath = os.path.relpath(filename, folder_paths.base_path)
|
||||
model_filenames = [
|
||||
os.path.abspath(line)
|
||||
for line in stdout.decode().splitlines()
|
||||
if not os.path.basename(line).startswith(".")
|
||||
]
|
||||
model_hashes = await async_batch_get_sha256(
|
||||
model_filenames,
|
||||
cache_only=not (ensure_sha or store_models),
|
||||
)
|
||||
|
||||
relpath_path = Path(relpath)
|
||||
for filename in model_filenames:
|
||||
relpath = os.path.relpath(filename, folder_paths.base_path)
|
||||
|
||||
model_data = {
|
||||
"filename": relpath,
|
||||
"sha256": _get_hash(Path(filename)),
|
||||
"explicit": relpath_path.name in used_inputs,
|
||||
"size": os.path.getsize(filename),
|
||||
"atime": os.path.getatime(filename),
|
||||
"ctime": os.path.getctime(filename),
|
||||
"disabled": relpath not in model_filter
|
||||
if model_filter is not None
|
||||
else False,
|
||||
"sha256": model_hashes.get(filename),
|
||||
}
|
||||
if store_models:
|
||||
|
||||
model_data["source"] = await alookup_model_source(
|
||||
model_data["sha256"],
|
||||
cache_only=not ensure_source,
|
||||
)
|
||||
# should_store = store_models and (
|
||||
# model_data["source"].get("source") != "huggingface"
|
||||
# or model_data["source"].get("repo", "").startswith("datasets/")
|
||||
# ) # TODO: sort this out
|
||||
should_store = store_models
|
||||
|
||||
if should_store:
|
||||
import bentoml
|
||||
|
||||
model_tag = f'cpack-model:{model_data["sha256"][:16]}'
|
||||
@@ -135,12 +163,14 @@ async def _get_models(data: dict, store_models: bool = False) -> list:
|
||||
model = bentoml.models.get(model_tag)
|
||||
except bentoml.exceptions.NotFound:
|
||||
with bentoml.models.create(
|
||||
model_tag,
|
||||
labels={"filename": filename},
|
||||
model_tag, labels={"filename": relpath}
|
||||
) as model:
|
||||
shutil.copy(filename, model.path_of("model.bin"))
|
||||
model_data["model_tag"] = model_tag
|
||||
models.append(model_data)
|
||||
if workflow_api:
|
||||
for model in models:
|
||||
model["refered"] = _is_file_refered(Path(model["filename"]), workflow_api)
|
||||
return models
|
||||
|
||||
|
||||
@@ -174,6 +204,7 @@ async def _get_custom_nodes() -> list:
|
||||
"url": url,
|
||||
"commit_hash": commit_hash,
|
||||
"disabled": subdir.name.endswith(".disabled"),
|
||||
"path": str(subdir.relative_to(custom_nodes)),
|
||||
}
|
||||
|
||||
for subdir in Path(custom_nodes).iterdir():
|
||||
@@ -199,15 +230,16 @@ async def _write_inputs(path: ZPath, data: dict) -> None:
|
||||
|
||||
input_dir = folder_paths.get_input_directory()
|
||||
|
||||
used_inputs = set()
|
||||
for node in data["workflow_api"].values():
|
||||
for _, v in node["inputs"].items():
|
||||
if isinstance(v, str):
|
||||
used_inputs.add(v)
|
||||
if "files" in data:
|
||||
selected = "\n".join(set(data.get("files", [])))
|
||||
else:
|
||||
selected = None
|
||||
|
||||
src_root = Path(input_dir).absolute()
|
||||
for src in src_root.glob("**/*"):
|
||||
rel = src.relative_to(src_root)
|
||||
if selected is not None and str(rel) not in selected:
|
||||
continue
|
||||
if src.is_dir():
|
||||
if isinstance(path, Path):
|
||||
path.joinpath("input").joinpath(rel).mkdir(parents=True, exist_ok=True)
|
||||
@@ -230,22 +262,261 @@ async def pack_workspace(request):
|
||||
|
||||
with zipfile.ZipFile(TEMP_FOLDER / zip_filename, "w") as zf:
|
||||
path = zipfile.Path(zf)
|
||||
await _write_requirements(path)
|
||||
await _write_snapshot(path, data)
|
||||
await _write_workflow(path, data)
|
||||
await _write_inputs(path, data)
|
||||
await _prepare_pack(path, data)
|
||||
|
||||
return web.json_response({"download_url": f"/bentoml/download/{zip_filename}"})
|
||||
|
||||
|
||||
class DevServer:
|
||||
TIMEOUT = 3600 * 24
|
||||
proc: Union[None, subprocess.Popen] = None
|
||||
watch_dog_task: asyncio.Task | None = None
|
||||
last_feed = 0
|
||||
run_dir: Path | None = None
|
||||
port = 0
|
||||
|
||||
@classmethod
|
||||
def start(cls, workflow_api: dict, port: int = 3000):
|
||||
cls.stop()
|
||||
|
||||
cls.port = port
|
||||
# prepare a temporary directory
|
||||
cls.run_dir = Path(tempfile.mkdtemp(suffix="-bento", prefix="comfy-pack-"))
|
||||
with cls.run_dir.joinpath("workflow_api.json").open("w") as f:
|
||||
f.write(json.dumps(workflow_api, indent=2))
|
||||
shutil.copy(
|
||||
Path(__file__).with_name("service.py"),
|
||||
cls.run_dir / "service.py",
|
||||
)
|
||||
shutil.copytree(COMFY_PACK_DIR, cls.run_dir / COMFY_PACK_DIR.name)
|
||||
|
||||
# find a free port
|
||||
self_port = 8188
|
||||
for i, arg in enumerate(sys.argv):
|
||||
if arg == "--port" or arg == "-p":
|
||||
self_port = int(sys.argv[i + 1])
|
||||
break
|
||||
|
||||
print(f"Starting dev server at port {port}, comfyui at port {self_port}")
|
||||
cls.proc = subprocess.Popen(
|
||||
[
|
||||
sys.executable,
|
||||
"-m",
|
||||
"bentoml",
|
||||
"serve",
|
||||
"service:ComfyService",
|
||||
"--port",
|
||||
str(port),
|
||||
],
|
||||
cwd=str(cls.run_dir.absolute()),
|
||||
env={
|
||||
**os.environ,
|
||||
"COMFYUI_SERVER": f"localhost:{self_port}",
|
||||
},
|
||||
)
|
||||
cls.watch_dog_task = asyncio.create_task(cls.watch_dog())
|
||||
cls.last_feed = time.time()
|
||||
|
||||
@classmethod
|
||||
async def watch_dog(cls):
|
||||
while True:
|
||||
await asyncio.sleep(0.1)
|
||||
if cls.last_feed + cls.TIMEOUT < time.time():
|
||||
cls.stop()
|
||||
break
|
||||
|
||||
@classmethod
|
||||
def feed_watch_dog(cls):
|
||||
if cls.proc:
|
||||
if cls.proc.poll() is None:
|
||||
cls.last_feed = time.time()
|
||||
return True
|
||||
else:
|
||||
cls.stop()
|
||||
return False
|
||||
return False
|
||||
|
||||
@classmethod
|
||||
def stop(cls):
|
||||
if cls.proc:
|
||||
cls.proc.terminate()
|
||||
cls.proc.wait()
|
||||
cls.proc = None
|
||||
time.sleep(1)
|
||||
print("Dev server stopped")
|
||||
if cls.watch_dog_task:
|
||||
cls.watch_dog_task.cancel()
|
||||
cls.watch_dog_task = None
|
||||
cls.last_feed = 0
|
||||
if cls.run_dir:
|
||||
shutil.rmtree(cls.run_dir)
|
||||
cls.run_dir = None
|
||||
|
||||
|
||||
def _parse_workflow(workflow: dict) -> tuple[dict[str, Any], dict[str, Any]]:
|
||||
inputs = {}
|
||||
outputs = {}
|
||||
dep_map = {}
|
||||
|
||||
for id, node in workflow.items():
|
||||
for input_name, v in node["inputs"].items():
|
||||
if isinstance(v, list) and len(v) == 2: # is a link
|
||||
dep_map[tuple(v)] = node, input_name
|
||||
|
||||
for id, node in workflow.items():
|
||||
node["id"] = id
|
||||
if node["class_type"].startswith("CPackInput"):
|
||||
if not node.get("inputs"):
|
||||
continue
|
||||
inputs[id] = node
|
||||
elif node["class_type"].startswith("CPackOutput"):
|
||||
if not node.get("inputs"):
|
||||
continue
|
||||
outputs[id] = node
|
||||
|
||||
return inputs, outputs
|
||||
|
||||
|
||||
def _validate_workflow(data: dict):
|
||||
workflow = data.get("workflow_api", {})
|
||||
if not workflow:
|
||||
return web.json_response(
|
||||
{
|
||||
"result": "error",
|
||||
"error": "empty workflow",
|
||||
},
|
||||
)
|
||||
input_spec, output_spec = _parse_workflow(workflow)
|
||||
if not input_spec:
|
||||
return web.json_response(
|
||||
{
|
||||
"result": "error",
|
||||
"error": "At least one ComfyPack input node is required",
|
||||
},
|
||||
)
|
||||
if not output_spec:
|
||||
return web.json_response(
|
||||
{
|
||||
"result": "error",
|
||||
"error": "At least one ComfyPack output node is required",
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
@PromptServer.instance.routes.post("/bentoml/serve")
|
||||
async def serve(request):
|
||||
data = await request.json()
|
||||
|
||||
if (error := _validate_workflow(data)) is not None:
|
||||
return error
|
||||
|
||||
DevServer.stop()
|
||||
|
||||
if _is_port_in_use(data.get("port", 3000), host=data.get("host", "localhost")):
|
||||
return web.json_response(
|
||||
{
|
||||
"result": "error",
|
||||
"error": "Port is already in use",
|
||||
},
|
||||
)
|
||||
try:
|
||||
DevServer.start(workflow_api=data["workflow_api"], port=data.get("port", 3000))
|
||||
return web.json_response(
|
||||
{
|
||||
"result": "success",
|
||||
"url": f"http://{data.get('host', 'localhost')}:{data.get('port', 3000)}",
|
||||
},
|
||||
)
|
||||
except Exception as e:
|
||||
return web.json_response(
|
||||
{
|
||||
"result": "error",
|
||||
"error": f"Build failed: {e.__class__.__name__}: {e}",
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
@PromptServer.instance.routes.post("/bentoml/serve/heartbeat")
|
||||
async def heartbeat(_):
|
||||
running = DevServer.feed_watch_dog()
|
||||
|
||||
if running:
|
||||
return web.json_response({"ready": True})
|
||||
else:
|
||||
return web.json_response({"error": "Server is not running"})
|
||||
|
||||
|
||||
@PromptServer.instance.routes.post("/bentoml/serve/terminate")
|
||||
async def terminate(_):
|
||||
DevServer.stop()
|
||||
return web.json_response({"result": "success"})
|
||||
|
||||
|
||||
@PromptServer.instance.routes.get("/bentoml/download/{zip_filename}")
|
||||
async def download_workspace(request):
|
||||
zip_filename = request.match_info["zip_filename"]
|
||||
return web.FileResponse(TEMP_FOLDER / zip_filename)
|
||||
|
||||
|
||||
async def _prepare_pack(
|
||||
working_dir: ZPath,
|
||||
data: dict,
|
||||
store_models: bool = False,
|
||||
) -> None:
|
||||
model_filter = set(data.get("models", []))
|
||||
models = await _get_models(
|
||||
store_models=store_models,
|
||||
model_filter=model_filter,
|
||||
)
|
||||
|
||||
await _write_requirements(working_dir, ["comfy-cli", "fastapi"])
|
||||
await _write_snapshot(working_dir, data, models)
|
||||
await _write_workflow(working_dir, data)
|
||||
await _write_inputs(working_dir, data)
|
||||
|
||||
|
||||
@PromptServer.instance.routes.post("/bentoml/model/query")
|
||||
async def get_models(request):
|
||||
data = await request.json()
|
||||
models = await _get_models(
|
||||
workflow_api=data.get("workflow_api"),
|
||||
ensure_sha=False,
|
||||
ensure_source=False,
|
||||
)
|
||||
return web.json_response({"models": models})
|
||||
|
||||
|
||||
async def _get_inputs(workflow_api):
|
||||
input_dir = folder_paths.get_input_directory()
|
||||
inputs = []
|
||||
for src in Path(input_dir).rglob("*"):
|
||||
if src.is_file():
|
||||
rel = src.relative_to(input_dir)
|
||||
badges = []
|
||||
checked = False
|
||||
if _is_file_refered(src, workflow_api):
|
||||
badges.append({"text": "Referenced"})
|
||||
checked = True
|
||||
data = {
|
||||
"path": str(rel),
|
||||
"badges": badges,
|
||||
"checked": checked,
|
||||
}
|
||||
inputs.append(data)
|
||||
return inputs
|
||||
|
||||
|
||||
@PromptServer.instance.routes.post("/bentoml/file/query")
|
||||
async def get_inputs(request):
|
||||
data = await request.json()
|
||||
inputs = await _get_inputs(
|
||||
workflow_api=data.get("workflow_api"),
|
||||
)
|
||||
return web.json_response({"files": inputs})
|
||||
|
||||
|
||||
@PromptServer.instance.routes.post("/bentoml/build")
|
||||
async def build_bento(request):
|
||||
async def build_bento_api(request):
|
||||
"""Request body: {
|
||||
workflow_api: dict,
|
||||
workflow: dict,
|
||||
@@ -259,42 +530,19 @@ async def build_bento(request):
|
||||
|
||||
data = await request.json()
|
||||
|
||||
if (error := _validate_workflow(data)) is not None:
|
||||
return error
|
||||
|
||||
with tempfile.TemporaryDirectory(suffix="-bento", prefix="comfy-pack-") as temp_dir:
|
||||
temp_dir_path = Path(temp_dir)
|
||||
# copy comfy_pack source code into the bento
|
||||
shutil.copytree(COMFY_PACK_DIR, temp_dir_path / COMFY_PACK_DIR.name)
|
||||
models = await _get_models(data, store_models=True)
|
||||
|
||||
await _write_requirements(temp_dir_path, ["comfy-cli", "fastapi"])
|
||||
await _write_snapshot(temp_dir_path, data, models)
|
||||
await _write_workflow(temp_dir_path, data)
|
||||
await _write_inputs(temp_dir_path, data)
|
||||
shutil.copy(
|
||||
Path(__file__).with_name("service.py"),
|
||||
temp_dir_path / "service.py",
|
||||
)
|
||||
await _prepare_pack(temp_dir_path, data, store_models=True)
|
||||
|
||||
# create a bento
|
||||
try:
|
||||
bento = bentoml.build(
|
||||
"service:ComfyService",
|
||||
name=data["bento_name"],
|
||||
build_ctx=temp_dir,
|
||||
models=[m["model_tag"] for m in models if "model_tag" in m],
|
||||
docker={
|
||||
"python_version": f"{sys.version_info.major}.{sys.version_info.minor}",
|
||||
"system_packages": [
|
||||
"git",
|
||||
"libglib2.0-0",
|
||||
"libsm6",
|
||||
"libxrender1",
|
||||
"libxext6",
|
||||
"ffmpeg",
|
||||
"libstdc++-12-dev",
|
||||
*data.get("system_packages", []),
|
||||
],
|
||||
},
|
||||
python={"requirements_txt": "requirements.txt", "lock_packages": True},
|
||||
bento = build_bento(
|
||||
data["bento_name"],
|
||||
temp_dir_path,
|
||||
data.get("system_packages"),
|
||||
)
|
||||
except bentoml.exceptions.BentoMLException as e:
|
||||
return web.json_response(
|
||||
|
||||
@@ -1,121 +0,0 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import shutil
|
||||
from functools import lru_cache
|
||||
from pathlib import Path
|
||||
from typing import Any, cast
|
||||
|
||||
import bentoml
|
||||
import fastapi
|
||||
|
||||
import comfy_pack
|
||||
import comfy_pack.run
|
||||
|
||||
REQUEST_TIMEOUT = 360
|
||||
BASE_DIR = Path(__file__).parent
|
||||
WORKFLOW_FILE = BASE_DIR / "workflow_api.json"
|
||||
COPY_THRESHOLD = 10 * 1024 * 1024
|
||||
INPUT_DIR = BASE_DIR / "input"
|
||||
logger = logging.getLogger("bentoml.service")
|
||||
|
||||
with open(WORKFLOW_FILE, "r") as f:
|
||||
workflow = json.load(f)
|
||||
|
||||
InputModel = comfy_pack.generate_input_model(workflow)
|
||||
app = fastapi.FastAPI()
|
||||
|
||||
|
||||
@lru_cache
|
||||
def _get_workspace() -> Path:
|
||||
import hashlib
|
||||
|
||||
from bentoml._internal.configuration.containers import BentoMLContainer
|
||||
|
||||
snapshot = BASE_DIR / "snapshot.json"
|
||||
checksum = hashlib.md5(snapshot.read_bytes()).hexdigest()
|
||||
wp = Path(BentoMLContainer.bentoml_home.get()) / "run/comfy_workspace" / checksum
|
||||
wp.parent.mkdir(parents=True, exist_ok=True)
|
||||
return wp
|
||||
|
||||
|
||||
@app.get("/workflow.json")
|
||||
def workflow_json():
|
||||
return workflow
|
||||
|
||||
|
||||
@bentoml.mount_asgi_app(app, path="/comfy")
|
||||
@bentoml.service(traffic={"timeout": REQUEST_TIMEOUT * 2})
|
||||
class ComfyService:
|
||||
def __init__(self):
|
||||
self.comfy_proc = comfy_pack.run.WorkflowRunner(
|
||||
str(_get_workspace()),
|
||||
str(INPUT_DIR),
|
||||
)
|
||||
self.comfy_proc.start(verbose=int("BENTOML_DEBUG" in os.environ))
|
||||
|
||||
@bentoml.api(input_spec=InputModel)
|
||||
def generate(
|
||||
self,
|
||||
*,
|
||||
ctx: bentoml.Context,
|
||||
**kwargs: Any,
|
||||
) -> Path:
|
||||
verbose = int("BENTOML_DEBUG" in os.environ)
|
||||
ret = self.comfy_proc.run_workflow(
|
||||
workflow,
|
||||
output_dir=ctx.temp_dir,
|
||||
timeout=REQUEST_TIMEOUT,
|
||||
verbose=verbose,
|
||||
**kwargs,
|
||||
)
|
||||
if isinstance(ret, list):
|
||||
ret = ret[-1]
|
||||
return ret
|
||||
|
||||
@bentoml.on_shutdown
|
||||
def on_shutdown(self):
|
||||
self.comfy_proc.stop()
|
||||
|
||||
@bentoml.on_deployment
|
||||
@staticmethod
|
||||
def prepare_comfy_workspace():
|
||||
from comfy_pack.package import install_comfyui, install_custom_modules
|
||||
|
||||
verbose = int("BENTOML_DEBUG" in os.environ)
|
||||
comfy_workspace = _get_workspace()
|
||||
|
||||
with BASE_DIR.joinpath("snapshot.json").open("rb") as f:
|
||||
snapshot = json.load(f)
|
||||
|
||||
if not comfy_workspace.joinpath(".DONE").exists():
|
||||
if comfy_workspace.exists():
|
||||
logger.info("Removing existing workspace")
|
||||
shutil.rmtree(comfy_workspace, ignore_errors=True)
|
||||
install_comfyui(snapshot, comfy_workspace, verbose=verbose)
|
||||
|
||||
for model in snapshot["models"]:
|
||||
model_tag = model.get("model_tag")
|
||||
if not model_tag:
|
||||
logger.warning(
|
||||
"Model %s is not in model store, the workflow may not work",
|
||||
model["filename"],
|
||||
)
|
||||
continue
|
||||
model_path = comfy_workspace / cast(str, model["filename"])
|
||||
model_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
bento_model = bentoml.models.get(model_tag)
|
||||
model_file = bento_model.path_of("model.bin")
|
||||
logger.info("Copying %s to %s", model_file, model_path)
|
||||
model_path.symlink_to(model_file)
|
||||
|
||||
for f in INPUT_DIR.glob("*"):
|
||||
if f.is_file():
|
||||
shutil.copy(f, comfy_workspace / "input" / f.name)
|
||||
elif f.is_dir():
|
||||
shutil.copytree(f, comfy_workspace / "input" / f.name)
|
||||
|
||||
install_custom_modules(snapshot, comfy_workspace, verbose=verbose)
|
||||
comfy_workspace.joinpath(".DONE").touch()
|
||||
+5
-3
@@ -1,6 +1,9 @@
|
||||
[project]
|
||||
name = "comfy-pack"
|
||||
description = "ComfyUI Interface Definition Language"
|
||||
description = """\
|
||||
A comprehensive toolkit for standardizing, packaging and deploying ComfyUI workflows \
|
||||
as reproducible environments and production-ready REST services\
|
||||
"""
|
||||
authors = [{ name = "Frost Ming", email = "frost@bentoml.com" }]
|
||||
readme = "README.md"
|
||||
requires-python = ">=3.9"
|
||||
@@ -8,7 +11,6 @@ dependencies = [
|
||||
"bentoml>=1.3.13",
|
||||
"click>=8.1.7",
|
||||
"comfy-cli>=1.2.8",
|
||||
"pydantic>=2.9",
|
||||
]
|
||||
dynamic = ["version"]
|
||||
|
||||
@@ -19,7 +21,7 @@ classifiers = [
|
||||
]
|
||||
|
||||
[project.urls]
|
||||
Homepage = "https://github.com/bentoml/ComfyUI-IDL"
|
||||
Homepage = "https://github.com/bentoml/comfy-pack"
|
||||
|
||||
[project.scripts]
|
||||
comfy-pack = "comfy_pack.cli:main"
|
||||
|
||||
@@ -1,2 +1,4 @@
|
||||
bentoml
|
||||
fastapi
|
||||
comfy-cli
|
||||
duckduckgo-search
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
from .run import WorkflowRunner
|
||||
from .run import ComfyUIServer, run_workflow
|
||||
from .utils import (
|
||||
generate_input_model,
|
||||
parse_workflow,
|
||||
@@ -7,7 +7,8 @@ from .utils import (
|
||||
)
|
||||
|
||||
__all__ = [
|
||||
"WorkflowRunner",
|
||||
"ComfyUIServer",
|
||||
"run_workflow",
|
||||
"parse_workflow",
|
||||
"generate_input_model",
|
||||
"populate_workflow",
|
||||
|
||||
+311
-53
@@ -1,28 +1,244 @@
|
||||
import click
|
||||
import hashlib
|
||||
import functools
|
||||
import json
|
||||
from pathlib import Path
|
||||
import os
|
||||
import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
from pathlib import Path
|
||||
|
||||
import click
|
||||
|
||||
from .const import COMFY_PACK_REPO, COMFYUI_MANAGER_REPO, COMFYUI_REPO, WORKSPACE_DIR
|
||||
from .hash import get_sha256
|
||||
from .utils import get_self_git_commit
|
||||
|
||||
|
||||
def _ensure_uv() -> None:
|
||||
"""Ensure uv is installed, raise error if not."""
|
||||
try:
|
||||
subprocess.run(
|
||||
["uv", "--version"],
|
||||
check=True,
|
||||
capture_output=True,
|
||||
)
|
||||
except (subprocess.SubprocessError, FileNotFoundError):
|
||||
raise RuntimeError(
|
||||
"uv is not installed. Please install it first:\n"
|
||||
"curl -LsSf https://astral.sh/uv/install.sh | sh"
|
||||
)
|
||||
|
||||
|
||||
@click.group()
|
||||
def main():
|
||||
"""comfy-pack CLI"""
|
||||
pass
|
||||
|
||||
|
||||
@main.command(
|
||||
name="init",
|
||||
help="Install latest ComfyUI and comfy-pack custom nodes and create a virtual environment",
|
||||
)
|
||||
@click.option(
|
||||
"--dir",
|
||||
"-d",
|
||||
default="ComfyUI",
|
||||
help="Target directory to install ComfyUI",
|
||||
type=click.Path(file_okay=False),
|
||||
)
|
||||
@click.option(
|
||||
"--verbose",
|
||||
"-v",
|
||||
count=True,
|
||||
help="Increase verbosity level (use multiple times for more verbosity)",
|
||||
help="Increase verbosity level",
|
||||
)
|
||||
@click.pass_context
|
||||
def main(ctx, verbose):
|
||||
"""ComfyUI Pack CLI"""
|
||||
ctx.ensure_object(dict)
|
||||
ctx.obj["verbose"] = verbose
|
||||
def init(dir: str, verbose: int):
|
||||
import os
|
||||
|
||||
from rich.console import Console
|
||||
|
||||
console = Console()
|
||||
|
||||
# Check if directory path is valid
|
||||
try:
|
||||
install_dir = Path(dir).absolute()
|
||||
if install_dir.exists() and not install_dir.is_dir():
|
||||
console.print(f"[red]Error: {dir} exists but is not a directory[/red]")
|
||||
return 1
|
||||
|
||||
# Check if directory is empty or contains ComfyUI
|
||||
if install_dir.exists():
|
||||
contents = list(install_dir.iterdir())
|
||||
if contents and not (install_dir / ".git").exists():
|
||||
console.print(
|
||||
f"[red]Error: Directory {dir} is not empty and doesn't appear to be a ComfyUI installation[/red]"
|
||||
)
|
||||
return 1
|
||||
except Exception as e:
|
||||
console.print(f"[red]Error: Invalid directory path - {str(e)}[/red]")
|
||||
return 1
|
||||
|
||||
# Check git installation
|
||||
try:
|
||||
subprocess.run(
|
||||
["git", "--version"],
|
||||
check=True,
|
||||
capture_output=True,
|
||||
)
|
||||
except (subprocess.SubprocessError, FileNotFoundError):
|
||||
console.print("[red]Error: git is not installed or not in PATH[/red]")
|
||||
return 1
|
||||
|
||||
# Check if we have write permissions
|
||||
try:
|
||||
if not install_dir.exists():
|
||||
install_dir.mkdir(parents=True)
|
||||
test_file = install_dir / ".write_test"
|
||||
test_file.touch()
|
||||
test_file.unlink()
|
||||
except (OSError, PermissionError) as e:
|
||||
console.print(f"[red]Error: No write permission in {dir} - {str(e)}[/red]")
|
||||
return 1
|
||||
|
||||
# Check if Python version is compatible
|
||||
if sys.version_info < (3, 8):
|
||||
console.print("[red]Error: Python 3.8 or higher is required[/red]")
|
||||
return 1
|
||||
|
||||
# Check if uv is installed
|
||||
try:
|
||||
_ensure_uv()
|
||||
except RuntimeError as e:
|
||||
console.print(f"[red]Error: {str(e)}[/red]")
|
||||
return 1
|
||||
|
||||
# Check if enough disk space is available (rough estimate: 2GB)
|
||||
try:
|
||||
free_space = shutil.disk_usage(install_dir).free
|
||||
if free_space < 2 * 1024 * 1024 * 1024: # 2GB in bytes
|
||||
console.print(
|
||||
"[yellow]Warning: Less than 2GB free disk space available[/yellow]"
|
||||
)
|
||||
except Exception as e:
|
||||
console.print(
|
||||
f"[yellow]Warning: Could not check free disk space - {str(e)}[/yellow]"
|
||||
)
|
||||
|
||||
# Clone ComfyUI if not exists
|
||||
if not (install_dir / ".git").exists():
|
||||
console.print("[green]Cloning ComfyUI...[/green]")
|
||||
subprocess.run(
|
||||
[
|
||||
"git",
|
||||
"clone",
|
||||
COMFYUI_REPO,
|
||||
str(install_dir),
|
||||
],
|
||||
check=True,
|
||||
)
|
||||
|
||||
# Update ComfyUI
|
||||
console.print("[green]Updating ComfyUI...[/green]")
|
||||
subprocess.run(
|
||||
["git", "pull"],
|
||||
cwd=install_dir,
|
||||
check=True,
|
||||
)
|
||||
|
||||
# Create and activate venv
|
||||
venv_dir = install_dir / ".venv"
|
||||
console.print("[green]Creating virtual environment with uv...[/green]")
|
||||
if venv_dir.exists():
|
||||
shutil.rmtree(venv_dir)
|
||||
subprocess.run(
|
||||
["uv", "venv", str(venv_dir)],
|
||||
check=True,
|
||||
)
|
||||
|
||||
# Get python path for future use
|
||||
if sys.platform == "win32":
|
||||
python = str(venv_dir / "Scripts" / "python.exe")
|
||||
|
||||
else:
|
||||
python = str(venv_dir / "bin" / "python")
|
||||
|
||||
# Install requirements with uv
|
||||
console.print("[green]Installing ComfyUI requirements with uv...[/green]")
|
||||
subprocess.run(
|
||||
["uv", "pip", "install", "pip", "--upgrade"],
|
||||
env={
|
||||
"VIRTUAL_ENV": str(venv_dir),
|
||||
"PATH": str(venv_dir / "bin") + os.pathsep + os.environ["PATH"],
|
||||
},
|
||||
check=True,
|
||||
)
|
||||
subprocess.run(
|
||||
["uv", "pip", "install", "-r", str(install_dir / "requirements.txt")],
|
||||
env={
|
||||
"VIRTUAL_ENV": str(venv_dir),
|
||||
"PATH": str(venv_dir / "bin") + os.pathsep + os.environ["PATH"],
|
||||
},
|
||||
check=True,
|
||||
)
|
||||
|
||||
# Install comfy-pack as custom node
|
||||
console.print("[green]Installing comfy-pack custom nodes...[/green]")
|
||||
custom_nodes_dir = install_dir / "custom_nodes"
|
||||
custom_nodes_dir.mkdir(exist_ok=True)
|
||||
|
||||
comfyui_manager_dir = custom_nodes_dir / "ComfyUI-Manager"
|
||||
if not (comfyui_manager_dir / ".git").exists():
|
||||
# Clone ComfyUI-Manager
|
||||
subprocess.run(
|
||||
["git", "clone", COMFYUI_MANAGER_REPO, str(comfyui_manager_dir)],
|
||||
check=True,
|
||||
)
|
||||
|
||||
comfy_pack_dir = custom_nodes_dir / "comfy-pack"
|
||||
if not (comfy_pack_dir / ".git").exists():
|
||||
# Clone comfy-pack
|
||||
subprocess.run(
|
||||
["git", "clone", COMFY_PACK_REPO, str(comfy_pack_dir)],
|
||||
check=True,
|
||||
)
|
||||
|
||||
# Update comfy-pack
|
||||
subprocess.run(
|
||||
["git", "pull"],
|
||||
cwd=comfy_pack_dir,
|
||||
check=True,
|
||||
)
|
||||
|
||||
# Install comfy-pack requirements
|
||||
if (comfy_pack_dir / "requirements.txt").exists():
|
||||
subprocess.run(
|
||||
[
|
||||
python,
|
||||
"-m",
|
||||
"pip",
|
||||
"install",
|
||||
"-r",
|
||||
str(comfy_pack_dir / "requirements.txt"),
|
||||
],
|
||||
check=True,
|
||||
)
|
||||
|
||||
version = get_self_git_commit() or "unknown"
|
||||
console.print(
|
||||
f"\n[green]✓ Installation completed! (comfy-pack version: {version})[/green]"
|
||||
)
|
||||
console.print(f"ComfyUI directory: {install_dir}")
|
||||
|
||||
console.print(
|
||||
"\n[green]Next steps:[/green]\n"
|
||||
f"1. cd {dir}\n"
|
||||
"2. source .venv/bin/activate # On Windows: .venv\\Scripts\\activate\n"
|
||||
"3. python main.py"
|
||||
)
|
||||
|
||||
|
||||
@main.command(
|
||||
name="restore",
|
||||
name="unpack",
|
||||
help="Restore the ComfyUI workspace to specified directory",
|
||||
)
|
||||
@click.argument("cpack", type=click.Path(exists=True, dir_okay=False))
|
||||
@@ -33,14 +249,26 @@ def main(ctx, verbose):
|
||||
help="target directory to restore the ComfyUI project",
|
||||
type=click.Path(file_okay=False),
|
||||
)
|
||||
@click.pass_context
|
||||
def restore_cmd(ctx, cpack: str, dir: str):
|
||||
from .package import install
|
||||
@click.option(
|
||||
"--include-disabled-models",
|
||||
default=False,
|
||||
type=click.BOOL,
|
||||
is_flag=True,
|
||||
)
|
||||
@click.option(
|
||||
"--verbose",
|
||||
"-v",
|
||||
count=True,
|
||||
help="Increase verbosity level (use multiple times for more verbosity)",
|
||||
)
|
||||
def unpack_cmd(cpack: str, dir: str, include_disabled_models: bool, verbose: int):
|
||||
from rich.console import Console
|
||||
|
||||
from .package import install
|
||||
|
||||
console = Console()
|
||||
|
||||
install(cpack, dir, ctx.obj["verbose"])
|
||||
install(cpack, dir, verbose=verbose, all_models=include_disabled_models)
|
||||
console.print("\n[green]✓ ComfyUI Workspace is restored![/green]")
|
||||
console.print(f"{dir}")
|
||||
|
||||
@@ -53,8 +281,8 @@ def restore_cmd(ctx, cpack: str, dir: str):
|
||||
|
||||
|
||||
def _print_schema(schema, verbose: int = 0):
|
||||
from rich.table import Table
|
||||
from rich.console import Console
|
||||
from rich.table import Table
|
||||
|
||||
table = Table(title="")
|
||||
|
||||
@@ -87,33 +315,10 @@ def _print_schema(schema, verbose: int = 0):
|
||||
Console().print(table)
|
||||
|
||||
|
||||
@main.command(name="info")
|
||||
@click.argument("cpack", type=click.Path(exists=True, dir_okay=False))
|
||||
@click.pass_context
|
||||
def info_cmd(ctx, cpack: str):
|
||||
"""
|
||||
Display information about the ComfyUI package.
|
||||
|
||||
Example:
|
||||
$ comfy_pack info workspace.cpack.zip
|
||||
"""
|
||||
from .utils import generate_input_model
|
||||
|
||||
with tempfile.TemporaryDirectory() as temp_dir:
|
||||
pack_dir = Path(temp_dir) / ".cpack"
|
||||
shutil.unpack_archive(cpack, pack_dir)
|
||||
workflow = json.loads((pack_dir / "workflow_api.json").read_text())
|
||||
|
||||
inputs = generate_input_model(workflow)
|
||||
_print_schema(inputs.model_json_schema(), ctx.obj["verbose"])
|
||||
|
||||
|
||||
@functools.lru_cache
|
||||
def _get_cache_workspace(cpack: str):
|
||||
m = hashlib.sha256()
|
||||
with open(cpack, "rb") as f:
|
||||
m.update(f.read())
|
||||
return Path.home() / ".comfypack" / "workspace" / m.hexdigest()[0:8]
|
||||
sha = get_sha256(cpack)
|
||||
return WORKSPACE_DIR / sha[0:8]
|
||||
|
||||
|
||||
@main.command(
|
||||
@@ -122,15 +327,24 @@ def _get_cache_workspace(cpack: str):
|
||||
"allow_extra_args": True,
|
||||
},
|
||||
help="Run a ComfyUI package with the given inputs",
|
||||
add_help_option=False,
|
||||
)
|
||||
@click.argument("cpack", type=click.Path(exists=True, dir_okay=False))
|
||||
@click.option("--output-dir", "-o", type=click.Path(), default=".")
|
||||
@click.option("--help", "-h", is_flag=True, help="Show this message and input schema")
|
||||
@click.option(
|
||||
"--verbose",
|
||||
"-v",
|
||||
count=True,
|
||||
help="Increase verbosity level (use multiple times for more verbosity)",
|
||||
)
|
||||
@click.pass_context
|
||||
def run(ctx, cpack: str, output_dir: str):
|
||||
from .utils import generate_input_model
|
||||
def run(ctx, cpack: str, output_dir: str, help: bool, verbose: int):
|
||||
from pydantic import ValidationError
|
||||
from rich.console import Console
|
||||
|
||||
from .utils import generate_input_model
|
||||
|
||||
inputs = dict(
|
||||
zip([k.lstrip("-").replace("-", "_") for k in ctx.args[::2]], ctx.args[1::2])
|
||||
)
|
||||
@@ -144,6 +358,15 @@ def run(ctx, cpack: str, output_dir: str):
|
||||
|
||||
input_model = generate_input_model(workflow)
|
||||
|
||||
# If help is requested, show command help and input schema
|
||||
if help:
|
||||
console.print(
|
||||
'Usage: comfy-pack run [OPTIONS] CPACK --input1 "value1" --input2 "value2" ...'
|
||||
)
|
||||
console.print("Run a ComfyUI package with the given inputs:")
|
||||
_print_schema(input_model.model_json_schema(), verbose)
|
||||
return 0
|
||||
|
||||
try:
|
||||
validated_data = input_model(**inputs)
|
||||
console.print("[green]✓ Input is valid![/green]")
|
||||
@@ -153,6 +376,9 @@ def run(ctx, cpack: str, output_dir: str):
|
||||
console.print("[red]✗ Validation failed![/red]")
|
||||
for error in e.errors():
|
||||
console.print(f"- {error['loc'][0]}: {error['msg']}")
|
||||
|
||||
console.print("\n[yellow]Expected inputs:[/yellow]")
|
||||
_print_schema(input_model.model_json_schema(), verbose)
|
||||
return 1
|
||||
|
||||
from .package import install
|
||||
@@ -162,23 +388,22 @@ def run(ctx, cpack: str, output_dir: str):
|
||||
console.print("\n[green]✓ Restoring ComfyUI Workspace...[/green]")
|
||||
if workspace.exists():
|
||||
shutil.rmtree(workspace)
|
||||
install(cpack, workspace)
|
||||
install(cpack, workspace, verbose=verbose)
|
||||
with open(workspace / "DONE", "w") as f:
|
||||
f.write("DONE")
|
||||
console.print("\n[green]✓ ComfyUI Workspace is restored![/green]")
|
||||
console.print(f"{workspace}")
|
||||
|
||||
from .run import WorkflowRunner
|
||||
from .run import ComfyUIServer, run_workflow
|
||||
|
||||
runner = None
|
||||
try:
|
||||
runner = WorkflowRunner(str(workspace.absolute()))
|
||||
runner.start()
|
||||
with ComfyUIServer(str(workspace.absolute()), verbose=verbose) as server:
|
||||
console.print("\n[green]✓ ComfyUI is launched in the background![/green]")
|
||||
results = runner.run_workflow(
|
||||
results = run_workflow(
|
||||
server.host,
|
||||
server.port,
|
||||
workflow,
|
||||
Path(output_dir).absolute(),
|
||||
verbose=ctx.obj["verbose"],
|
||||
verbose=verbose,
|
||||
**validated_data.model_dump(),
|
||||
)
|
||||
console.print("\n[green]✓ Workflow is executed successfully![/green]")
|
||||
@@ -192,6 +417,39 @@ def run(ctx, cpack: str, output_dir: str):
|
||||
console.print(f"{i}: {value}")
|
||||
else:
|
||||
console.print(results)
|
||||
finally:
|
||||
if runner:
|
||||
runner.stop(verbose=ctx.obj["verbose"])
|
||||
|
||||
|
||||
@main.command(name="build-bento")
|
||||
@click.argument("source")
|
||||
@click.option("--name", help="Name of the bento service")
|
||||
@click.option("--version", help="Version of the bento service")
|
||||
def bento_cmd(source: str, name: str | None, version: str | None):
|
||||
"""Build a bento from the source, which can be either a .cpack.zip file or a bento tag."""
|
||||
import bentoml
|
||||
from bentoml.bentos import BentoBuildConfig
|
||||
|
||||
from .package import build_bento
|
||||
|
||||
with tempfile.TemporaryDirectory() as temp_dir:
|
||||
if source.endswith(".cpack.zip"):
|
||||
name = name or os.path.basename(source).replace(".cpack.zip", "")
|
||||
shutil.unpack_archive(source, temp_dir)
|
||||
system_packages = None
|
||||
include_default_system_packages = True
|
||||
else:
|
||||
existing_bento = bentoml.get(source)
|
||||
name = name or existing_bento.tag.name
|
||||
shutil.copytree(existing_bento.path_of("src"), temp_dir, dirs_exist_ok=True)
|
||||
build_config = BentoBuildConfig.from_bento_dir(
|
||||
existing_bento.path_of("src")
|
||||
)
|
||||
system_packages = build_config.docker.system_packages
|
||||
include_default_system_packages = False
|
||||
|
||||
build_bento(
|
||||
name,
|
||||
Path(temp_dir),
|
||||
version=version,
|
||||
system_packages=system_packages,
|
||||
include_default_system_packages=include_default_system_packages,
|
||||
)
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
import pathlib
|
||||
import os
|
||||
|
||||
|
||||
CPACK_HOME = (
|
||||
pathlib.Path.home() / ".comfypack"
|
||||
if not os.environ.get("CPACK_HOME", "")
|
||||
else pathlib.Path(os.environ.get("CPACK_HOME", ""))
|
||||
)
|
||||
if not CPACK_HOME.exists():
|
||||
CPACK_HOME.mkdir()
|
||||
|
||||
MODEL_DIR = CPACK_HOME / "models"
|
||||
WORKSPACE_DIR = CPACK_HOME / "workspace"
|
||||
SHA_CACHE_FILE = CPACK_HOME / "sha_cache.json"
|
||||
MODEL_SOURCE_CACHE_FILE = CPACK_HOME / "model_source_cache.json"
|
||||
|
||||
COMFYUI_REPO = "https://github.com/comfyanonymous/ComfyUI.git"
|
||||
COMFY_PACK_REPO = "https://github.com/bentoml/comfy-pack.git"
|
||||
COMFYUI_MANAGER_REPO = "https://github.com/ltdrdata/ComfyUI-Manager.git"
|
||||
|
||||
STRICT_MODE = os.environ.get("CPACK_STRICT_MODE", "0") in ["1", "true", "True"]
|
||||
@@ -0,0 +1,122 @@
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import subprocess
|
||||
import sys
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
from datetime import datetime
|
||||
from functools import partial
|
||||
from typing import Dict, List
|
||||
|
||||
from .const import SHA_CACHE_FILE
|
||||
|
||||
CALC_CMD = """
|
||||
import hashlib
|
||||
import sys
|
||||
|
||||
filepath = sys.argv[1]
|
||||
chunk_size = int(sys.argv[2])
|
||||
|
||||
sha256 = hashlib.sha256()
|
||||
with open(filepath, "rb") as f:
|
||||
for chunk in iter(lambda: f.read(chunk_size), b""):
|
||||
sha256.update(chunk)
|
||||
print(sha256.hexdigest())
|
||||
"""
|
||||
|
||||
|
||||
def calculate_sha256_worker(filepath: str, chunk_size: int = 4 * 1024 * 1024) -> str:
|
||||
"""Calculate SHA-256 in a separate process"""
|
||||
result = subprocess.run(
|
||||
[sys.executable, "-c", CALC_CMD, filepath, str(chunk_size)],
|
||||
stdout=subprocess.PIPE,
|
||||
stderr=subprocess.PIPE,
|
||||
text=True,
|
||||
)
|
||||
assert result.returncode == 0, result.stderr
|
||||
return result.stdout.strip()
|
||||
|
||||
|
||||
def get_sha256(filepath: str) -> str:
|
||||
return batch_get_sha256([filepath])[filepath]
|
||||
|
||||
|
||||
def async_get_sha256(filepath: str) -> str:
|
||||
return asyncio.run(async_batch_get_sha256([filepath]))[filepath]
|
||||
|
||||
|
||||
def batch_get_sha256(filepaths: List[str], cache_only: bool = False) -> Dict[str, str]:
|
||||
return asyncio.run(async_batch_get_sha256(filepaths, cache_only=cache_only))
|
||||
|
||||
|
||||
async def async_batch_get_sha256(
|
||||
filepaths: List[str],
|
||||
cache_only: bool = False,
|
||||
) -> Dict[str, str]:
|
||||
# Load cache
|
||||
cache = {}
|
||||
if SHA_CACHE_FILE.exists():
|
||||
try:
|
||||
with SHA_CACHE_FILE.open("r") as f:
|
||||
cache = json.load(f)
|
||||
except (json.JSONDecodeError, IOError):
|
||||
pass
|
||||
|
||||
# Initialize process pool
|
||||
max_workers = max(1, (os.cpu_count() or 1))
|
||||
|
||||
# Process files
|
||||
results = {}
|
||||
new_cache = {}
|
||||
async with asyncio.Lock():
|
||||
with ThreadPoolExecutor(max_workers=max_workers) as pool:
|
||||
loop = asyncio.get_event_loop()
|
||||
|
||||
for filepath in filepaths:
|
||||
if not os.path.exists(filepath):
|
||||
results[filepath] = None
|
||||
continue
|
||||
|
||||
# Get file info
|
||||
stat = os.stat(filepath)
|
||||
current_size = stat.st_size
|
||||
current_time = stat.st_ctime
|
||||
|
||||
# Check cache
|
||||
cache_entry = cache.get(filepath)
|
||||
if cache_entry:
|
||||
if (
|
||||
cache_entry["size"] == current_size
|
||||
and cache_entry["birthtime"] == current_time
|
||||
):
|
||||
results[filepath] = cache_entry["sha256"]
|
||||
continue
|
||||
|
||||
if cache_only:
|
||||
results[filepath] = ""
|
||||
continue
|
||||
|
||||
# Calculate new SHA
|
||||
calc_func = partial(calculate_sha256_worker, filepath)
|
||||
sha256 = await loop.run_in_executor(pool, calc_func)
|
||||
|
||||
# Update cache and results
|
||||
new_cache[filepath] = {
|
||||
"sha256": sha256,
|
||||
"size": current_size,
|
||||
"birthtime": current_time,
|
||||
"last_verified": datetime.now().isoformat(),
|
||||
}
|
||||
results[filepath] = sha256
|
||||
|
||||
# Save cache
|
||||
try:
|
||||
with SHA_CACHE_FILE.open("r") as f:
|
||||
cache = json.load(f)
|
||||
cache.update(new_cache)
|
||||
with SHA_CACHE_FILE.open("w") as f:
|
||||
json.dump(cache, f, indent=2)
|
||||
except (IOError, OSError):
|
||||
pass
|
||||
|
||||
return results
|
||||
@@ -0,0 +1,170 @@
|
||||
from .const import MODEL_SOURCE_CACHE_FILE
|
||||
import asyncio
|
||||
import json
|
||||
import re
|
||||
|
||||
# MODEL_NAME = r"[a-zA-Z0-9-._]+"
|
||||
# COMMIT = r"[a-f0-9]+"
|
||||
|
||||
COMMIT_PATTERN = re.compile(r'href="/([a-zA-Z0-9-._/]+)/commit/([a-f0-9]+)"')
|
||||
|
||||
PATH_PATTERN = re.compile(
|
||||
r'data-target="CopyButton" data-props="{"value":"([^&]+)"'
|
||||
)
|
||||
|
||||
|
||||
async def _lookup_huggingface_model(model_sha: str) -> dict:
|
||||
from duckduckgo_search import DDGS
|
||||
import aiohttp
|
||||
|
||||
query = f"site:huggingface.co blob {model_sha}"
|
||||
|
||||
try:
|
||||
with DDGS() as ddgs:
|
||||
search_results = ddgs.text(query, max_results=5)
|
||||
|
||||
async with aiohttp.ClientSession(trust_env=True) as session:
|
||||
for result in search_results:
|
||||
url = result['link']
|
||||
if "blob" not in url:
|
||||
continue
|
||||
|
||||
try:
|
||||
async with session.get(url) as resp:
|
||||
if resp.status != 200:
|
||||
continue
|
||||
text = await resp.text()
|
||||
if commit_match := COMMIT_PATTERN.search(text):
|
||||
repo, commit = commit_match.groups()
|
||||
if path_match := PATH_PATTERN.search(text):
|
||||
path = path_match.group(1)
|
||||
download_url = f"https://huggingface.co/{repo}/resolve/{commit}/{path}?download=true"
|
||||
url = f"https://huggingface.co/{repo}/blob/{commit}/{path}"
|
||||
info = {
|
||||
"download_url": download_url,
|
||||
"url": url,
|
||||
"repo": repo,
|
||||
"commit": commit,
|
||||
"path": path,
|
||||
"source": "huggingface",
|
||||
}
|
||||
return info
|
||||
except aiohttp.ClientError:
|
||||
continue
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
return {}
|
||||
|
||||
|
||||
async def _loopup_civitai_model(model_sha: str) -> dict:
|
||||
import aiohttp
|
||||
|
||||
async with aiohttp.ClientSession(trust_env=True) as session:
|
||||
async with session.post(
|
||||
"https://meilisearch-v1-9.civitai.com/multi-search",
|
||||
headers={
|
||||
"accept": "*/*",
|
||||
"accept-language": "en,zh;q=0.9,zh-CN;q=0.8",
|
||||
"cache-control": "no-cache",
|
||||
"content-type": "application/json",
|
||||
"origin": "https://civitai.com",
|
||||
"pragma": "no-cache",
|
||||
"priority": "u=1, i",
|
||||
"referer": "https://civitai.com/",
|
||||
"sec-ch-ua": '"Google Chrome";v="131", "Chromium";v="131", "Not_A Brand";v="24"',
|
||||
"sec-ch-ua-mobile": "?0",
|
||||
"sec-ch-ua-platform": '"macOS"',
|
||||
"sec-fetch-dest": "empty",
|
||||
"sec-fetch-mode": "cors",
|
||||
"sec-fetch-site": "same-site",
|
||||
"user-agent": "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36",
|
||||
"x-meilisearch-client": "Meilisearch instant-meilisearch (v0.13.5) ; Meilisearch JavaScript (v0.34.0)",
|
||||
},
|
||||
json={
|
||||
"queries": [
|
||||
{
|
||||
"q": model_sha,
|
||||
"indexUid": "models_v9",
|
||||
"facets": [
|
||||
"category.name",
|
||||
"checkpointType",
|
||||
"fileFormats",
|
||||
"lastVersionAtUnix",
|
||||
"tags.name",
|
||||
"type",
|
||||
"user.username",
|
||||
"version.baseModel",
|
||||
],
|
||||
"attributesToHighlight": [],
|
||||
"highlightPreTag": "__ais-highlight__",
|
||||
"highlightPostTag": "__/ais-highlight__",
|
||||
"limit": 51,
|
||||
"offset": 0,
|
||||
"filter": ["nsfwLevel=1"],
|
||||
}
|
||||
]
|
||||
},
|
||||
) as resp:
|
||||
if resp.status != 200:
|
||||
return {}
|
||||
data = await resp.json()
|
||||
|
||||
if len(data.get("results", [])) == 0:
|
||||
return {}
|
||||
|
||||
if len(data["results"][0]["hits"]) == 0:
|
||||
return {}
|
||||
|
||||
hit = data["results"][0]["hits"][0]
|
||||
repo_id = hit["id"]
|
||||
repo_name = hit["name"]
|
||||
versions = hit["versions"]
|
||||
|
||||
for version in versions:
|
||||
if model_sha.upper() in version["hashes"]:
|
||||
break
|
||||
else:
|
||||
return {}
|
||||
version_id = version["id"]
|
||||
version_name = version["name"]
|
||||
return {
|
||||
"download_url": f"https://civitai.com/api/download/models/{version_id}",
|
||||
"url": f"https://civitai.com/models/{repo_id}?modelVersionId={version_id}",
|
||||
"repo": repo_id,
|
||||
"commit": version_id,
|
||||
"source": "civitai",
|
||||
"repo_name": repo_name,
|
||||
"version_name": version_name,
|
||||
}
|
||||
|
||||
|
||||
async def alookup_model_source(model_sha: str, cache_only=False) -> dict:
|
||||
if not model_sha:
|
||||
return {}
|
||||
try:
|
||||
model_source_cache = json.loads(MODEL_SOURCE_CACHE_FILE.read_text())
|
||||
except Exception:
|
||||
with open(MODEL_SOURCE_CACHE_FILE, "w") as f:
|
||||
f.write("{}")
|
||||
model_source_cache = {}
|
||||
if model_source_cache.get(model_sha):
|
||||
return model_source_cache[model_sha]
|
||||
if cache_only:
|
||||
return {}
|
||||
|
||||
info = await _lookup_huggingface_model(model_sha)
|
||||
if not info:
|
||||
info = await _loopup_civitai_model(model_sha)
|
||||
|
||||
# elemental read and write
|
||||
model_source_cache = json.loads(MODEL_SOURCE_CACHE_FILE.read_text())
|
||||
model_source_cache[model_sha] = info
|
||||
with open(MODEL_SOURCE_CACHE_FILE, "w") as f:
|
||||
json.dump(model_source_cache, f)
|
||||
|
||||
return info
|
||||
|
||||
|
||||
def lookup_model_source(model_sha: str, cache_only=False) -> dict:
|
||||
return asyncio.run(alookup_model_source(model_sha, cache_only=cache_only))
|
||||
+320
-16
@@ -6,9 +6,20 @@ import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
import threading
|
||||
import urllib.parse
|
||||
import urllib.request
|
||||
from pathlib import Path
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
COMFYUI_REPO = "https://github.com/comfyanonymous/ComfyUI.git"
|
||||
from .const import COMFYUI_REPO, MODEL_DIR, STRICT_MODE
|
||||
from .hash import get_sha256
|
||||
from .utils import get_self_git_commit
|
||||
|
||||
if TYPE_CHECKING:
|
||||
import bentoml
|
||||
|
||||
COMFY_PACK_DIR = Path(__file__).parent
|
||||
|
||||
|
||||
def _clone_commit(url: str, commit: str, dir: Path, verbose: int = 0):
|
||||
@@ -58,6 +69,9 @@ def install_custom_modules(snapshot, workspace: Path, verbose: int = 0):
|
||||
print("Installing custom nodes")
|
||||
for module in snapshot["custom_nodes"]:
|
||||
url = module["url"]
|
||||
if not url.strip():
|
||||
print(f"Skipping invalid custom node: {module}")
|
||||
continue
|
||||
directory = url.split("/")[-1].split(".")[0]
|
||||
module_dir = workspace / "custom_nodes" / directory
|
||||
|
||||
@@ -73,6 +87,7 @@ def install_custom_modules(snapshot, workspace: Path, verbose: int = 0):
|
||||
_clone_commit(url, commit_hash, module_dir, verbose=verbose)
|
||||
|
||||
if module_dir.joinpath("install.py").exists():
|
||||
env = os.environ.copy()
|
||||
venv = workspace / ".venv"
|
||||
if venv.exists():
|
||||
python = (
|
||||
@@ -80,8 +95,17 @@ def install_custom_modules(snapshot, workspace: Path, verbose: int = 0):
|
||||
if os.name == "nt"
|
||||
else venv / "bin" / "python"
|
||||
)
|
||||
if "PATH" in env:
|
||||
env["PATH"] = f"{str(python.parent)}:{env['PATH']}"
|
||||
else:
|
||||
env["PATH"] = str(python.parent)
|
||||
env["VIRTUAL_ENV"] = str(venv)
|
||||
else:
|
||||
python = Path(sys.executable)
|
||||
|
||||
if verbose > 0:
|
||||
print(f"Installing {directory} custom node")
|
||||
print(f"$ {python.absolute()} install.py")
|
||||
subprocess.check_call(
|
||||
[str(python.absolute()), "install.py"],
|
||||
cwd=module_dir,
|
||||
@@ -98,8 +122,7 @@ def install_dependencies(
|
||||
workspace: Path,
|
||||
verbose: int = 0,
|
||||
):
|
||||
if verbose > 0:
|
||||
print("Installing Python dependencies")
|
||||
print("Installing Python dependencies")
|
||||
python_version = snapshot["python"]
|
||||
stdout = None if verbose > 0 else subprocess.DEVNULL
|
||||
stderr = None if verbose > 1 else subprocess.DEVNULL
|
||||
@@ -139,17 +162,24 @@ def install_dependencies(
|
||||
stdout=stdout,
|
||||
stderr=stderr,
|
||||
)
|
||||
if verbose > 0:
|
||||
print(f"Installing dependencies from {req_file}")
|
||||
install_cmd = [
|
||||
"uv",
|
||||
"pip",
|
||||
"install",
|
||||
"-p",
|
||||
str(venv_py),
|
||||
"-r",
|
||||
req_file,
|
||||
"--no-deps",
|
||||
]
|
||||
if STRICT_MODE:
|
||||
pass
|
||||
else:
|
||||
install_cmd.extend(["--index-strategy", "unsafe-best-match"])
|
||||
subprocess.check_call(
|
||||
[
|
||||
"uv",
|
||||
"pip",
|
||||
"install",
|
||||
"-p",
|
||||
str(venv_py),
|
||||
"-r",
|
||||
req_file,
|
||||
"--no-deps",
|
||||
],
|
||||
install_cmd,
|
||||
stdout=stdout,
|
||||
stderr=stderr,
|
||||
)
|
||||
@@ -158,8 +188,206 @@ def install_dependencies(
|
||||
return venv_py
|
||||
|
||||
|
||||
def install(cpack: str | Path, workspace: str | Path = "workspace", verbose: int = 0):
|
||||
def get_search_url(sha: str) -> str:
|
||||
"""Generate custom search URLs for model on HuggingFace and CivitAI"""
|
||||
base_url = "https://duckduckgo.com"
|
||||
sha = sha.upper()
|
||||
hf_query = f"{sha} OR {sha[:10]}"
|
||||
hf_query = urllib.parse.quote(hf_query)
|
||||
return f"{base_url}?q={hf_query}"
|
||||
|
||||
|
||||
def download_file(url: str, dest_path: Path, progress_callback=None):
|
||||
"""Download file with progress tracking"""
|
||||
if subprocess.call(["curl", "--version"], stdout=subprocess.DEVNULL) == 0:
|
||||
subprocess.check_call(
|
||||
["curl", "-L", url, "-o", str(dest_path)],
|
||||
)
|
||||
return True
|
||||
try:
|
||||
with urllib.request.urlopen(url) as response:
|
||||
total_size = int(response.headers.get("content-length", 0))
|
||||
block_size = 8192
|
||||
downloaded = 0
|
||||
|
||||
with open(dest_path, "wb") as f:
|
||||
while True:
|
||||
buffer = response.read(block_size)
|
||||
if not buffer:
|
||||
break
|
||||
downloaded += len(buffer)
|
||||
f.write(buffer)
|
||||
if progress_callback:
|
||||
progress = (
|
||||
(downloaded / total_size) * 100 if total_size > 0 else 0
|
||||
)
|
||||
progress_callback(progress)
|
||||
return True
|
||||
except Exception as e:
|
||||
print(f"Download failed: {e}")
|
||||
if dest_path.exists():
|
||||
dest_path.unlink()
|
||||
return False
|
||||
|
||||
|
||||
def show_progress(filename: str):
|
||||
"""Progress callback function"""
|
||||
|
||||
def callback(progress):
|
||||
print(f"\rDownloading {filename}: {progress:.1f}%", end="")
|
||||
|
||||
return callback
|
||||
|
||||
|
||||
def create_model_symlink(global_path: Path, sha: str, target_path: Path, filename: str):
|
||||
"""Create symlink from global storage to workspace"""
|
||||
source = global_path / sha
|
||||
target = target_path / filename
|
||||
|
||||
if target.exists():
|
||||
if target.is_symlink():
|
||||
target.unlink()
|
||||
else:
|
||||
raise RuntimeError(f"File {target} already exists and is not a symlink")
|
||||
|
||||
target.parent.mkdir(parents=True, exist_ok=True)
|
||||
os.symlink(source, target)
|
||||
|
||||
|
||||
def retrieve_models(
|
||||
snapshot: dict,
|
||||
workspace: Path,
|
||||
download: bool = True,
|
||||
all_models: bool = False,
|
||||
verbose: int = 0,
|
||||
):
|
||||
"""Retrieve models from user downloads"""
|
||||
print("Retrieving models")
|
||||
models = snapshot.get("models", [])
|
||||
if not models:
|
||||
return
|
||||
|
||||
MODEL_DIR.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
for model in models:
|
||||
sha = model["sha256"]
|
||||
filename = model["filename"]
|
||||
disabled = model.get("disabled", False)
|
||||
if (workspace / filename).exists():
|
||||
if not (MODEL_DIR / sha).exists() and (workspace / filename).is_file():
|
||||
shutil.move(workspace / filename, MODEL_DIR / sha)
|
||||
create_model_symlink(MODEL_DIR, sha, workspace, filename)
|
||||
continue
|
||||
|
||||
if (MODEL_DIR / sha).exists():
|
||||
print(f"Model {filename} already exists in cache")
|
||||
create_model_symlink(MODEL_DIR, sha, workspace, filename)
|
||||
continue
|
||||
|
||||
if disabled and not all_models:
|
||||
continue
|
||||
|
||||
if not download:
|
||||
continue
|
||||
|
||||
print(f"\nModel {filename} is never downloaded before")
|
||||
if source := model.get("source"):
|
||||
url = source["download_url"]
|
||||
target_path = MODEL_DIR / sha
|
||||
download_thread = threading.Thread(
|
||||
target=download_file,
|
||||
args=(url, target_path, show_progress(filename)),
|
||||
)
|
||||
download_thread.start()
|
||||
download_thread.join()
|
||||
|
||||
if not target_path.exists():
|
||||
print("\nDownload failed!")
|
||||
continue
|
||||
|
||||
print("\nDownload completed! Verifying SHA256...")
|
||||
terget_sha = get_sha256(str(target_path))
|
||||
if terget_sha != sha:
|
||||
print("SHA256 verification failed! File may be corrupted or incorrect.")
|
||||
target_path.unlink()
|
||||
else:
|
||||
continue
|
||||
|
||||
search_url = get_search_url(sha)
|
||||
print(f"Search URL: {search_url}")
|
||||
print(f"Path: {workspace / filename}")
|
||||
|
||||
while True:
|
||||
path = input("Enter path to downloaded file (or 'skip' to skip): ")
|
||||
if path.lower() == "skip":
|
||||
break
|
||||
|
||||
try:
|
||||
# Check if input is a URL
|
||||
if path.startswith(("http://", "https://")):
|
||||
url = path
|
||||
target_path = MODEL_DIR / sha
|
||||
|
||||
# Start download in a separate thread
|
||||
download_thread = threading.Thread(
|
||||
target=download_file,
|
||||
args=(url, target_path, show_progress(filename)),
|
||||
)
|
||||
download_thread.start()
|
||||
download_thread.join()
|
||||
|
||||
if not target_path.exists():
|
||||
print("\nDownload failed!")
|
||||
continue
|
||||
|
||||
print("\nDownload completed! Verifying SHA256...")
|
||||
terget_sha = get_sha256(str(target_path))
|
||||
if terget_sha != sha:
|
||||
print(
|
||||
"SHA256 verification failed! File may be corrupted or incorrect."
|
||||
)
|
||||
target_path.unlink()
|
||||
continue
|
||||
|
||||
print("SHA256 verification successful!")
|
||||
else:
|
||||
# Handle local file
|
||||
downloaded_path = Path(path)
|
||||
if not downloaded_path.exists():
|
||||
print("File does not exist!")
|
||||
continue
|
||||
|
||||
# Verify SHA256 before copying
|
||||
print("Verifying SHA256...")
|
||||
target_sha = get_sha256(str(downloaded_path))
|
||||
if target_sha != sha:
|
||||
print(
|
||||
f"Downloaded file SHA256 does not match expected SHA256: {target_sha} != {sha}"
|
||||
)
|
||||
continue
|
||||
|
||||
print("SHA256 verification successful!")
|
||||
# Copy to global storage
|
||||
shutil.copy2(downloaded_path, MODEL_DIR / sha)
|
||||
|
||||
# Create symlink
|
||||
create_model_symlink(MODEL_DIR, sha, workspace, filename)
|
||||
print(f"Model {filename} installed successfully")
|
||||
break
|
||||
except Exception as e:
|
||||
print(f"Error processing file: {e}")
|
||||
continue
|
||||
|
||||
|
||||
def install(
|
||||
cpack: str | Path,
|
||||
workspace: str | Path = "workspace",
|
||||
preheat: bool = True,
|
||||
all_models: bool = False,
|
||||
verbose: int = 0,
|
||||
):
|
||||
workspace = Path(workspace)
|
||||
print(f"Installing package {cpack} to {workspace} (verbose={verbose})")
|
||||
with tempfile.TemporaryDirectory() as temp_dir:
|
||||
pack_dir = Path(temp_dir) / ".cpack"
|
||||
shutil.unpack_archive(cpack, pack_dir)
|
||||
@@ -168,10 +396,86 @@ def install(cpack: str | Path, workspace: str | Path = "workspace", verbose: int
|
||||
|
||||
install_comfyui(snapshot, workspace, verbose=verbose)
|
||||
install_dependencies(snapshot, str(req_txt_file), workspace, verbose=verbose)
|
||||
install_custom_modules(snapshot, workspace, verbose=verbose)
|
||||
|
||||
for f in (pack_dir / "input").glob("*"):
|
||||
if f.is_file():
|
||||
shutil.copy(f, workspace / "input" / f.name)
|
||||
elif f.is_dir():
|
||||
shutil.copytree(f, workspace / "input" / f.name)
|
||||
shutil.copytree(f, workspace / "input" / f.name, dirs_exist_ok=True)
|
||||
|
||||
retrieve_models(
|
||||
snapshot,
|
||||
workspace,
|
||||
verbose=verbose,
|
||||
download=False,
|
||||
)
|
||||
|
||||
install_custom_modules(snapshot, workspace, verbose=verbose)
|
||||
if preheat:
|
||||
from .run import ComfyUIServer
|
||||
|
||||
with ComfyUIServer(
|
||||
str(workspace),
|
||||
verbose=verbose,
|
||||
venv=str(workspace / ".venv"),
|
||||
) as _:
|
||||
pass
|
||||
|
||||
retrieve_models(snapshot, workspace, verbose=verbose, all_models=all_models)
|
||||
|
||||
|
||||
required_files = ["snapshot.json", "requirements.txt"]
|
||||
|
||||
|
||||
def build_bento(
|
||||
bento_name: str,
|
||||
source_dir: Path,
|
||||
*,
|
||||
version: str | None = None,
|
||||
system_packages: list[str] | None = None,
|
||||
include_default_system_packages: bool = True,
|
||||
) -> bentoml.Bento:
|
||||
import bentoml
|
||||
|
||||
for f in required_files:
|
||||
if not (source_dir / f).exists():
|
||||
raise FileNotFoundError(f"Not a valid comfy-pack package: missing `{f}`")
|
||||
|
||||
if include_default_system_packages:
|
||||
system_packages = [
|
||||
"git",
|
||||
"libglib2.0-0",
|
||||
"libsm6",
|
||||
"libxrender1",
|
||||
"libxext6",
|
||||
"ffmpeg",
|
||||
"libstdc++-12-dev",
|
||||
*(system_packages or []),
|
||||
]
|
||||
else:
|
||||
system_packages = system_packages or []
|
||||
|
||||
shutil.copy2(Path(__file__).with_name("service.py"), source_dir / "service.py")
|
||||
# Copy comfy-pack package
|
||||
shutil.copytree(
|
||||
COMFY_PACK_DIR, source_dir / COMFY_PACK_DIR.name, dirs_exist_ok=True
|
||||
)
|
||||
snapshot = json.loads((source_dir / "snapshot.json").read_text())
|
||||
return bentoml.build(
|
||||
"service:ComfyService",
|
||||
name=bento_name,
|
||||
version=version,
|
||||
build_ctx=str(source_dir),
|
||||
labels={"comfy-pack-version": get_self_git_commit() or "unknown"},
|
||||
models=[
|
||||
m["model_tag"]
|
||||
for m in snapshot["models"]
|
||||
if "model_tag" in m and not m.get("disabled", False)
|
||||
],
|
||||
docker={
|
||||
"python_version": f"{sys.version_info.major}.{sys.version_info.minor}",
|
||||
"system_packages": system_packages,
|
||||
"setup_script": Path(__file__).with_name("setup_workspace.py").as_posix(),
|
||||
},
|
||||
python={"requirements_txt": "requirements.txt", "lock_packages": True},
|
||||
)
|
||||
|
||||
+180
-95
@@ -1,12 +1,15 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import random
|
||||
import copy
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import random
|
||||
import shutil
|
||||
import socket
|
||||
import subprocess
|
||||
import time
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from typing import Any, Union
|
||||
|
||||
@@ -24,28 +27,63 @@ def _probe_comfyui_server(port: int) -> None:
|
||||
req = request.Request(full_url)
|
||||
_ = request.urlopen(req)
|
||||
|
||||
full_url = f"http://127.0.0.1:{port}/api/object_info"
|
||||
req = request.Request(full_url)
|
||||
_ = request.urlopen(req)
|
||||
|
||||
class WorkflowRunner:
|
||||
|
||||
def _is_port_in_use(port: int | str, host="localhost"):
|
||||
if isinstance(port, str):
|
||||
port = int(port)
|
||||
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
|
||||
try:
|
||||
s.connect((host, port))
|
||||
return True
|
||||
except ConnectionRefusedError:
|
||||
return False
|
||||
except Exception:
|
||||
return True
|
||||
|
||||
|
||||
class ComfyUIServer:
|
||||
def __init__(
|
||||
self,
|
||||
workspace: str,
|
||||
input_dir: str | None = None,
|
||||
host: str = "localhost",
|
||||
port: int | None = None,
|
||||
venv: str | None = None,
|
||||
verbose: int = 0,
|
||||
) -> None:
|
||||
"""
|
||||
Initialize the WorkflowRunner.
|
||||
|
||||
Args:
|
||||
workspace (str): The workspace path for ComfyUI.
|
||||
workspace (str, optional): The workspace path for ComfyUI. If not specified, runner will try to connect to an existing ComfyUI server.
|
||||
input_dir (str, optional): The input directory for ComfyUI. Defaults to None.
|
||||
port (int, optional): The port number for ComfyUI. Defaults to None. If 8188 is in use, a random port will be chosen.
|
||||
"""
|
||||
self.workspace = workspace
|
||||
self.temp_dir = Path(workspace) / "cli_run" / "temp"
|
||||
self.output_dir = Path(workspace) / "cli_run" / "output"
|
||||
self.input_dir = input_dir
|
||||
self.is_running = False
|
||||
self.port = port if port else random.randint(58000, 58999)
|
||||
self.verbose = verbose
|
||||
self.host = host
|
||||
self.server_proc: subprocess.Popen | None = None
|
||||
|
||||
def start(self, verbose: int = 0) -> None:
|
||||
run_dir = (Path(workspace) / "cli_run").absolute()
|
||||
self.temp_dir = run_dir / "temp"
|
||||
self.output_dir = run_dir / "output"
|
||||
|
||||
self.temp_dir.mkdir(parents=True, exist_ok=True)
|
||||
self.output_dir.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
if port is None:
|
||||
if _is_port_in_use(8188):
|
||||
self.port = port if port else random.randint(58000, 58999)
|
||||
else:
|
||||
self.port = 8188
|
||||
else:
|
||||
self.port = port
|
||||
self.venv = venv
|
||||
|
||||
def start(self) -> None:
|
||||
"""
|
||||
Start the ComfyUI process.
|
||||
|
||||
@@ -58,30 +96,31 @@ class WorkflowRunner:
|
||||
Raises:
|
||||
RuntimeError: If ComfyUI is already running.
|
||||
"""
|
||||
if self.is_running:
|
||||
raise RuntimeError("ComfyUI Runner is already started")
|
||||
|
||||
logger.info(
|
||||
"Disable tracking from Comfy CLI, not for privacy concerns, but to workaround a bug"
|
||||
)
|
||||
self.temp_dir.mkdir(parents=True, exist_ok=True)
|
||||
self.output_dir.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
stdout = None if verbose > 0 else subprocess.DEVNULL
|
||||
env = os.environ.copy()
|
||||
if self.venv:
|
||||
env["VIRTUAL_ENV"] = self.venv
|
||||
if os.name == "nt":
|
||||
env["PATH"] = f"{self.venv}\\Scripts;{env.get('PATH', '')}"
|
||||
else:
|
||||
env["PATH"] = f"{self.venv}/bin:{env.get('PATH', '')}"
|
||||
|
||||
stdout = None if self.verbose > 0 else subprocess.DEVNULL
|
||||
command = ["comfy", "--skip-prompt", "tracking", "disable"]
|
||||
subprocess.run(command, check=True, stdout=stdout)
|
||||
subprocess.run(command, check=True, stdout=stdout, env=env)
|
||||
logger.info("Successfully disabled Comfy CLI tracking")
|
||||
|
||||
logger.info("Preparing directories required by ComfyUI...")
|
||||
|
||||
logger.info("Starting ComfyUI in the background...")
|
||||
command = [
|
||||
"comfy",
|
||||
"--workspace",
|
||||
self.workspace,
|
||||
"launch",
|
||||
"--background",
|
||||
"--",
|
||||
"python",
|
||||
"main.py",
|
||||
"--output-directory",
|
||||
self.output_dir,
|
||||
"--temp-directory",
|
||||
@@ -91,14 +130,36 @@ class WorkflowRunner:
|
||||
]
|
||||
if self.input_dir:
|
||||
command.extend(["--input-directory", self.input_dir])
|
||||
if subprocess.run(command, check=True, stdout=stdout):
|
||||
self.is_running = True
|
||||
|
||||
if self.host != "localhost":
|
||||
command.extend(["--listen", self.host])
|
||||
|
||||
def preexec_fn():
|
||||
os.setpgrp()
|
||||
|
||||
self.server_proc = subprocess.Popen(
|
||||
command,
|
||||
stdout=stdout,
|
||||
stderr=None,
|
||||
preexec_fn=preexec_fn,
|
||||
env=env,
|
||||
cwd=self.workspace,
|
||||
)
|
||||
|
||||
if _wait_for_startup(self.host, self.port):
|
||||
_probe_comfyui_server(self.port)
|
||||
logger.info("Successfully started ComfyUI in the background")
|
||||
else:
|
||||
logger.error("Failed to start ComfyUI in the background")
|
||||
|
||||
def stop(self, verbose: int = 0) -> None:
|
||||
def is_running(self) -> bool:
|
||||
if self.server_proc is None:
|
||||
return False
|
||||
if self.server_proc.poll() is not None:
|
||||
return False
|
||||
return True
|
||||
|
||||
def stop(self) -> None:
|
||||
"""
|
||||
Stop the ComfyUI process.
|
||||
|
||||
@@ -107,13 +168,13 @@ class WorkflowRunner:
|
||||
Raises:
|
||||
RuntimeError: If ComfyUI is not currently running.
|
||||
"""
|
||||
if not self.is_running:
|
||||
raise RuntimeError("ComfyUI Runner is not started yet")
|
||||
|
||||
if self.server_proc is None:
|
||||
raise RuntimeError("ComfyUI server is not started yet")
|
||||
proc = self.server_proc
|
||||
self.server_proc = None
|
||||
logger.info("Stopping ComfyUI...")
|
||||
command = ["comfy", "stop"]
|
||||
stdout = None if verbose > 0 else subprocess.DEVNULL
|
||||
subprocess.run(command, check=True, stdout=stdout)
|
||||
proc.terminate()
|
||||
proc.wait()
|
||||
logger.info("Successfully stopped ComfyUI")
|
||||
|
||||
logger.info("Cleaning up temporary directory...")
|
||||
@@ -121,79 +182,103 @@ class WorkflowRunner:
|
||||
shutil.rmtree(self.output_dir, ignore_errors=True)
|
||||
logger.info("Successfully cleaned up temporary directory")
|
||||
|
||||
self.is_running = False
|
||||
def __enter__(self):
|
||||
self.start()
|
||||
return self
|
||||
|
||||
def run_workflow(
|
||||
self,
|
||||
workflow: dict,
|
||||
output_dir: Union[str, Path, None] = None,
|
||||
timeout: int = 300,
|
||||
verbose: int = 0,
|
||||
**kwargs: Any,
|
||||
) -> Any:
|
||||
"""
|
||||
Run a ComfyUI workflow.
|
||||
def __exit__(self, exc_type, exc_val, exc_tb):
|
||||
self.stop()
|
||||
|
||||
This method executes a given workflow, populates it with input data,
|
||||
and retrieves the output.
|
||||
|
||||
Args:
|
||||
workflow (dict): The workflow to run.
|
||||
output_dir (Union[str, Path, None], optional): Temporary directory for the workflow. Defaults to None.
|
||||
timeout (int, optional): Timeout for the workflow execution in seconds. Defaults to 300.
|
||||
**kwargs: Additional keyword arguments for workflow population.
|
||||
def _wait_for_startup(host: str, port: int, timeout: int = 1800) -> bool:
|
||||
start_time = time.time()
|
||||
while time.time() - start_time < timeout:
|
||||
if _is_port_in_use(port, host):
|
||||
return True
|
||||
time.sleep(1)
|
||||
return False
|
||||
|
||||
Returns:
|
||||
Any: The output of the workflow.
|
||||
|
||||
Raises:
|
||||
RuntimeError: If ComfyUI is not started.
|
||||
"""
|
||||
if not self.is_running:
|
||||
raise RuntimeError("ComfyUI Runner is not started yet")
|
||||
def run_workflow(
|
||||
host: str,
|
||||
port: int,
|
||||
workflow: dict,
|
||||
output_dir: Union[str, Path, None] = None,
|
||||
timeout: int = 300,
|
||||
verbose: int = 0,
|
||||
**kwargs: Any,
|
||||
) -> Any:
|
||||
"""
|
||||
Run a ComfyUI workflow.
|
||||
|
||||
workflow_copy = copy.deepcopy(workflow)
|
||||
if output_dir is None:
|
||||
output_dir = Path(".")
|
||||
if isinstance(output_dir, str):
|
||||
output_dir = Path(output_dir)
|
||||
This method executes a given workflow, populates it with input data,
|
||||
and retrieves the output.
|
||||
|
||||
run_id = os.urandom(8).hex()
|
||||
populate_workflow(
|
||||
workflow_copy,
|
||||
output_dir,
|
||||
session_id=run_id,
|
||||
**kwargs,
|
||||
)
|
||||
Args:
|
||||
workflow (dict): The workflow to run.
|
||||
output_dir (Union[str, Path, None], optional): Temporary directory for the workflow. Defaults to None.
|
||||
timeout (int, optional): Timeout for the workflow execution in seconds. Defaults to 300.
|
||||
**kwargs: Additional keyword arguments for workflow population.
|
||||
|
||||
workflow_file_path = Path(self.workspace) / "workflow.json"
|
||||
with open(workflow_file_path, "w") as file:
|
||||
json.dump(workflow_copy, file)
|
||||
Returns:
|
||||
Any: The output of the workflow.
|
||||
|
||||
extra_args = []
|
||||
if verbose > 0:
|
||||
extra_args.append("--verbose")
|
||||
Raises:
|
||||
RuntimeError: If ComfyUI is not started.
|
||||
"""
|
||||
run_id = uuid.uuid4().hex[:8]
|
||||
|
||||
# Execute the workflow
|
||||
command = [
|
||||
"comfy",
|
||||
"run",
|
||||
"--workflow",
|
||||
workflow_file_path.as_posix(),
|
||||
"--port",
|
||||
str(self.port),
|
||||
"--timeout",
|
||||
str(timeout),
|
||||
"--wait",
|
||||
*extra_args,
|
||||
]
|
||||
env = os.environ.copy()
|
||||
env["NO_COLOR"] = "1"
|
||||
subprocess.run(command, check=True, env=env)
|
||||
workflow_copy = copy.deepcopy(workflow)
|
||||
if output_dir is None:
|
||||
output_dir = Path(".")
|
||||
if isinstance(output_dir, str):
|
||||
output_dir = Path(output_dir)
|
||||
|
||||
# retrieve the output
|
||||
return retrieve_workflow_outputs(
|
||||
workflow_copy,
|
||||
output_dir,
|
||||
session_id=run_id,
|
||||
)
|
||||
run_id = os.urandom(8).hex()
|
||||
populate_workflow(
|
||||
workflow_copy,
|
||||
output_dir,
|
||||
session_id=run_id,
|
||||
**kwargs,
|
||||
)
|
||||
|
||||
workflow_file_path = output_dir / f"workflow_{run_id}.json"
|
||||
with open(workflow_file_path, "w") as file:
|
||||
json.dump(workflow_copy, file)
|
||||
|
||||
extra_args = []
|
||||
if verbose > 0:
|
||||
extra_args.append("--verbose")
|
||||
|
||||
stdout = None if verbose > 0 else subprocess.DEVNULL
|
||||
command = ["comfy", "--skip-prompt", "tracking", "disable"]
|
||||
subprocess.run(command, check=True, stdout=stdout)
|
||||
|
||||
# Execute the workflow
|
||||
command = [
|
||||
"comfy",
|
||||
"--skip-prompt",
|
||||
"run",
|
||||
"--workflow",
|
||||
workflow_file_path.as_posix(),
|
||||
"--port",
|
||||
str(port),
|
||||
"--host",
|
||||
host,
|
||||
"--timeout",
|
||||
str(timeout),
|
||||
"--wait",
|
||||
*extra_args,
|
||||
]
|
||||
env = os.environ.copy()
|
||||
env["NO_COLOR"] = "1"
|
||||
subprocess.run(command, check=True, env=env)
|
||||
|
||||
workflow_file_path.unlink()
|
||||
|
||||
# retrieve the output
|
||||
return retrieve_workflow_outputs(
|
||||
workflow_copy,
|
||||
output_dir,
|
||||
session_id=run_id,
|
||||
)
|
||||
|
||||
@@ -0,0 +1,189 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import signal
|
||||
import threading
|
||||
import time
|
||||
from functools import lru_cache
|
||||
from pathlib import Path
|
||||
from typing import Any, cast
|
||||
|
||||
import bentoml
|
||||
import fastapi
|
||||
from bentoml.models import HuggingFaceModel
|
||||
|
||||
import comfy_pack
|
||||
import comfy_pack.run
|
||||
|
||||
REQUEST_TIMEOUT = 3600
|
||||
BASE_DIR = Path(__file__).parent
|
||||
COPY_THRESHOLD = 10 * 1024 * 1024
|
||||
INPUT_DIR = BASE_DIR / "input"
|
||||
logger = logging.getLogger("bentoml.service")
|
||||
|
||||
|
||||
EXISTING_COMFYUI_SERVER = os.environ.get("COMFYUI_SERVER")
|
||||
|
||||
with BASE_DIR.joinpath("workflow_api.json").open() as f:
|
||||
workflow = json.load(f)
|
||||
|
||||
InputModel = comfy_pack.generate_input_model(workflow)
|
||||
app = fastapi.FastAPI()
|
||||
|
||||
|
||||
@lru_cache
|
||||
def _get_workspace() -> Path:
|
||||
import hashlib
|
||||
|
||||
from bentoml._internal.configuration.containers import BentoMLContainer
|
||||
|
||||
snapshot = BASE_DIR / "snapshot.json"
|
||||
checksum = hashlib.md5(snapshot.read_bytes()).hexdigest()
|
||||
wp = (
|
||||
Path(BentoMLContainer.bentoml_home.get()) / "run" / "comfy_workspace" / checksum
|
||||
)
|
||||
wp.parent.mkdir(parents=True, exist_ok=True)
|
||||
return wp
|
||||
|
||||
|
||||
@app.get("/workflow.json")
|
||||
def workflow_json():
|
||||
return workflow
|
||||
|
||||
|
||||
def _watch_server(server: comfy_pack.run.ComfyUIServer):
|
||||
while True:
|
||||
time.sleep(1)
|
||||
if not server.is_running():
|
||||
if server.server_proc is not None:
|
||||
logger.warning(
|
||||
"Server exited with code %s", server.server_proc.returncode
|
||||
)
|
||||
os.kill(os.getpid(), signal.SIGTERM)
|
||||
break
|
||||
|
||||
|
||||
if not EXISTING_COMFYUI_SERVER:
|
||||
# register models
|
||||
with BASE_DIR.joinpath("snapshot.json").open("rb") as f:
|
||||
snapshot = json.load(f)
|
||||
else:
|
||||
snapshot = {}
|
||||
|
||||
|
||||
@bentoml.mount_asgi_app(app, path="/comfy")
|
||||
@bentoml.service(traffic={"timeout": REQUEST_TIMEOUT * 2}, resources={"gpu": 1})
|
||||
class ComfyService:
|
||||
def __init__(self):
|
||||
logger = logging.getLogger("comfy_pack")
|
||||
logger.setLevel(logging.INFO)
|
||||
if not EXISTING_COMFYUI_SERVER:
|
||||
self.server = comfy_pack.run.ComfyUIServer(
|
||||
str(_get_workspace()),
|
||||
str(INPUT_DIR),
|
||||
verbose=int("BENTOML_DEBUG" in os.environ),
|
||||
)
|
||||
self.server.start()
|
||||
logger.info(
|
||||
"ComfyUI Server started at %s:%s", self.server.host, self.server.port
|
||||
)
|
||||
self.host = self.server.host
|
||||
self.port = self.server.port
|
||||
self.watch_thread = threading.Thread(
|
||||
target=_watch_server,
|
||||
args=(self.server,),
|
||||
daemon=True,
|
||||
)
|
||||
self.watch_thread.start()
|
||||
logger.info("Watch thread started")
|
||||
else:
|
||||
logger.info("Attaching to ComfyUI server: %s", EXISTING_COMFYUI_SERVER)
|
||||
if ":" in EXISTING_COMFYUI_SERVER:
|
||||
self.host, port = EXISTING_COMFYUI_SERVER.split(":")
|
||||
self.port = int(port)
|
||||
else:
|
||||
self.host = EXISTING_COMFYUI_SERVER
|
||||
self.port = 80
|
||||
|
||||
@bentoml.api(input_spec=InputModel)
|
||||
def generate(
|
||||
self,
|
||||
*,
|
||||
ctx: bentoml.Context,
|
||||
**kwargs: Any,
|
||||
) -> Path:
|
||||
verbose = int("BENTOML_DEBUG" in os.environ)
|
||||
ret = comfy_pack.run_workflow(
|
||||
self.host,
|
||||
self.port,
|
||||
workflow,
|
||||
output_dir=ctx.temp_dir,
|
||||
timeout=REQUEST_TIMEOUT,
|
||||
verbose=verbose,
|
||||
**kwargs,
|
||||
)
|
||||
if isinstance(ret, list):
|
||||
ret = ret[-1]
|
||||
return ret
|
||||
|
||||
@bentoml.on_shutdown
|
||||
def on_shutdown(self):
|
||||
logger.info("Shutting down")
|
||||
if not EXISTING_COMFYUI_SERVER:
|
||||
self.watch_thread.join()
|
||||
logger.info("Watch thread finished")
|
||||
|
||||
@bentoml.on_deployment
|
||||
@staticmethod
|
||||
def prepare_models():
|
||||
if EXISTING_COMFYUI_SERVER:
|
||||
return
|
||||
comfy_workspace = _get_workspace()
|
||||
if not comfy_workspace.joinpath(".DONE").exists():
|
||||
raise RuntimeError("ComfyUI workspace is not ready")
|
||||
for model in snapshot["models"]:
|
||||
if model.get("disabled", False):
|
||||
continue
|
||||
model_path = comfy_workspace / cast(str, model["filename"])
|
||||
if model_path.exists():
|
||||
continue
|
||||
if model_tag := model.get("model_tag"):
|
||||
model_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
bento_model = bentoml.models.get(model_tag)
|
||||
model_file = bento_model.path_of("model.bin")
|
||||
print(f"Copying {model_file} to {model_path}")
|
||||
model_path.symlink_to(model_file)
|
||||
elif (source := model["source"]).get("source") == "huggingface":
|
||||
matched = next(
|
||||
(
|
||||
m
|
||||
for m in ComfyService.models
|
||||
if isinstance(m, HuggingFaceModel)
|
||||
and m.model_id.lower() == source["repo"].lower()
|
||||
and source["commit"].lower() == m.revision.lower()
|
||||
),
|
||||
None,
|
||||
)
|
||||
if matched is not None:
|
||||
model_file = os.path.join(matched.resolve(), source["path"])
|
||||
model_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
print(f"Copying {model_file} to {model_path}")
|
||||
model_path.symlink_to(model_file)
|
||||
else:
|
||||
print(
|
||||
f"WARN: Unrecognized model source: {source}, the model may be missing"
|
||||
)
|
||||
|
||||
|
||||
if False and not EXISTING_COMFYUI_SERVER:
|
||||
for model in snapshot["models"]:
|
||||
if model.get("disabled"):
|
||||
continue
|
||||
source = model["source"]
|
||||
if source.get("source") != "huggingface" or source["repo"].startswith(
|
||||
"datasets/"
|
||||
):
|
||||
continue
|
||||
ComfyService.models.append(HuggingFaceModel(source["repo"], source["commit"]))
|
||||
Executable
+65
@@ -0,0 +1,65 @@
|
||||
#!/usr/bin/env python3
|
||||
# ruff: noqa: E402
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import subprocess
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
virtualenv = os.environ.get("VIRTUAL_ENV")
|
||||
if virtualenv and "--reload" not in sys.argv:
|
||||
print("Re-executing in virtualenv:", virtualenv)
|
||||
venv_python = os.path.join(virtualenv, "bin/python3")
|
||||
os.execl(venv_python, venv_python, *sys.argv, "--reload")
|
||||
|
||||
# The script path is ./env/docker/setup_script
|
||||
SRC_DIR = Path(__file__).parent.parent.parent / "src"
|
||||
INPUT_DIR = SRC_DIR / "input"
|
||||
sys.path.append(str(SRC_DIR))
|
||||
|
||||
|
||||
def _get_workspace() -> tuple[Path, dict]:
|
||||
import hashlib
|
||||
import json
|
||||
|
||||
from bentoml._internal.configuration.containers import BentoMLContainer
|
||||
|
||||
snapshot = SRC_DIR / "snapshot.json"
|
||||
checksum = hashlib.md5(snapshot.read_bytes()).hexdigest()
|
||||
wp = (
|
||||
Path(BentoMLContainer.bentoml_home.get()) / "run" / "comfy_workspace" / checksum
|
||||
)
|
||||
wp.parent.mkdir(parents=True, exist_ok=True)
|
||||
return wp, json.loads(snapshot.read_text())
|
||||
|
||||
|
||||
def prepare_comfy_workspace():
|
||||
import shutil
|
||||
|
||||
from comfy_pack.package import install_comfyui, install_custom_modules
|
||||
|
||||
verbose = int("BENTOML_DEBUG" in os.environ)
|
||||
comfy_workspace, snapshot = _get_workspace()
|
||||
|
||||
if not comfy_workspace.joinpath(".DONE").exists():
|
||||
if comfy_workspace.exists():
|
||||
print("Removing existing workspace")
|
||||
shutil.rmtree(comfy_workspace, ignore_errors=True)
|
||||
install_comfyui(snapshot, comfy_workspace, verbose=verbose)
|
||||
|
||||
for f in INPUT_DIR.glob("*"):
|
||||
if f.is_file():
|
||||
shutil.copy(f, comfy_workspace / "input" / f.name)
|
||||
elif f.is_dir():
|
||||
shutil.copytree(f, comfy_workspace / "input" / f.name)
|
||||
|
||||
install_custom_modules(snapshot, comfy_workspace, verbose=verbose)
|
||||
comfy_workspace.joinpath(".DONE").touch()
|
||||
subprocess.run(
|
||||
["chown", "-R", "bentoml:bentoml", str(comfy_workspace)], check=True
|
||||
)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
prepare_comfy_workspace()
|
||||
+43
-1
@@ -1,8 +1,9 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import subprocess
|
||||
from pathlib import Path
|
||||
import re
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from typing import TYPE_CHECKING, Any, Literal, Union
|
||||
|
||||
if TYPE_CHECKING:
|
||||
@@ -258,3 +259,44 @@ def retrieve_workflow_outputs(
|
||||
if len(outs) == 1:
|
||||
return outs[0]
|
||||
return outs
|
||||
|
||||
|
||||
def get_self_git_commit() -> str | None:
|
||||
"""Get current git commit of the repository.
|
||||
|
||||
Returns:
|
||||
str | None: Git commit hash in format "{hash}[-dirty]" or None if not in a git repo
|
||||
"""
|
||||
try:
|
||||
repo_root = Path(__file__).parent.parent.parent
|
||||
|
||||
# Check if we're in a git repo
|
||||
subprocess.run(
|
||||
["git", "rev-parse", "--git-dir"],
|
||||
cwd=repo_root,
|
||||
check=True,
|
||||
capture_output=True,
|
||||
)
|
||||
|
||||
# Get current commit hash
|
||||
commit_hash = subprocess.run(
|
||||
["git", "rev-parse", "--short", "HEAD"],
|
||||
cwd=repo_root,
|
||||
check=True,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
).stdout.strip()
|
||||
|
||||
# Check if working directory is clean
|
||||
is_dirty = (
|
||||
subprocess.run(
|
||||
["git", "diff", "--quiet"],
|
||||
cwd=repo_root,
|
||||
check=False,
|
||||
).returncode
|
||||
!= 0
|
||||
)
|
||||
|
||||
return f"{commit_hash}-dirty" if is_dirty else commit_hash
|
||||
except (subprocess.SubprocessError, FileNotFoundError):
|
||||
return None
|
||||
|
||||
@@ -523,13 +523,12 @@ wheels = [
|
||||
|
||||
[[package]]
|
||||
name = "comfy-pack"
|
||||
version = "0.1.4.dev10+g1d19c05"
|
||||
version = "0.2.2.dev7+g5238115.d20241231"
|
||||
source = { editable = "." }
|
||||
dependencies = [
|
||||
{ name = "bentoml" },
|
||||
{ name = "click" },
|
||||
{ name = "comfy-cli" },
|
||||
{ name = "pydantic" },
|
||||
]
|
||||
|
||||
[package.metadata]
|
||||
@@ -537,7 +536,6 @@ requires-dist = [
|
||||
{ name = "bentoml", specifier = ">=1.3.13" },
|
||||
{ name = "click", specifier = ">=8.1.7" },
|
||||
{ name = "comfy-cli", specifier = ">=1.2.8" },
|
||||
{ name = "pydantic", specifier = ">=2.9" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
|
||||
+1067
-123
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user