diff --git a/VERSION b/VERSION index 45c7a584..22c08f72 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -v0.0.1 +v0.2.1 diff --git a/deploy/.gitignore b/deploy/.gitignore deleted file mode 100644 index 41f567e1..00000000 --- a/deploy/.gitignore +++ /dev/null @@ -1,5 +0,0 @@ -/docker-compose/mysql/data/ -/docker-compose/flexmodel/storage/ -/docker-compose/flexmodel/pages/ -/docker-compose/flexmodel/logs/ -/docker-compose/flexmodel/database/ \ No newline at end of file diff --git a/deploy/docker-compose/.env b/deploy/docker-compose/001-mysql-single/.env similarity index 100% rename from deploy/docker-compose/.env rename to deploy/docker-compose/001-mysql-single/.env diff --git a/deploy/docker-compose/001-mysql-single/.gitignore b/deploy/docker-compose/001-mysql-single/.gitignore new file mode 100644 index 00000000..dcf94be2 --- /dev/null +++ b/deploy/docker-compose/001-mysql-single/.gitignore @@ -0,0 +1,5 @@ +/mysql/data/ +/flexmodel/storage/ +/flexmodel/pages/ +/flexmodel/logs/ +/flexmodel/database/ \ No newline at end of file diff --git a/deploy/docker-compose/README.md b/deploy/docker-compose/001-mysql-single/README.md similarity index 100% rename from deploy/docker-compose/README.md rename to deploy/docker-compose/001-mysql-single/README.md diff --git a/deploy/docker-compose/docker-compose.yml b/deploy/docker-compose/001-mysql-single/docker-compose.yml similarity index 100% rename from deploy/docker-compose/docker-compose.yml rename to deploy/docker-compose/001-mysql-single/docker-compose.yml diff --git a/deploy/docker-compose/mysql/init/import.sql b/deploy/docker-compose/001-mysql-single/mysql/init/import.sql similarity index 100% rename from deploy/docker-compose/mysql/init/import.sql rename to deploy/docker-compose/001-mysql-single/mysql/init/import.sql diff --git a/deploy/docker-compose/mysql/my.cnf b/deploy/docker-compose/001-mysql-single/mysql/my.cnf similarity index 100% rename from deploy/docker-compose/mysql/my.cnf rename to deploy/docker-compose/001-mysql-single/mysql/my.cnf diff --git a/deploy/docker-compose/nginx/conf.d/default.conf.template b/deploy/docker-compose/001-mysql-single/nginx/conf.d/default.conf.template similarity index 100% rename from deploy/docker-compose/nginx/conf.d/default.conf.template rename to deploy/docker-compose/001-mysql-single/nginx/conf.d/default.conf.template diff --git a/deploy/docker-compose/nginx/nginx.conf b/deploy/docker-compose/001-mysql-single/nginx/nginx.conf similarity index 100% rename from deploy/docker-compose/nginx/nginx.conf rename to deploy/docker-compose/001-mysql-single/nginx/nginx.conf diff --git a/deploy/docker-compose/002-mysql-fullstack/.env b/deploy/docker-compose/002-mysql-fullstack/.env new file mode 100644 index 00000000..6e3f4424 --- /dev/null +++ b/deploy/docker-compose/002-mysql-fullstack/.env @@ -0,0 +1,35 @@ +CONTAINER_NAME_PREFIX=prod2 + +# 运行模式: native (GraalVM 原生镜像) 或 jvm (JVM 模式) +# 切换方式:修改 FLEXMODEL_RUN_MODE 和 FLEXMODEL_SERVER_IMAGE 两行 +# native → FLEXMODEL_SERVER_IMAGE=cjbi/flexmodel-server-native:latest +# jvm → FLEXMODEL_SERVER_IMAGE=cjbi/flexmodel-server:latest +FLEXMODEL_SERVER_IMAGE=cjbi/flexmodel-server-native:latest + +# mysql +MYSQL_DATABASE=flexmodel +MYSQL_ROOT_PASSWORD=123456 + +# nginx +NGINX_HTTP_PORT=80 +NGINX_HTTPS_PORT=443 + +# rabbitmq +RABBITMQ_DEFAULT_USER=flexmodel +RABBITMQ_DEFAULT_PASS=flexmodel +RABBITMQ_AMQP_PORT=5672 +RABBITMQ_MANAGEMENT_PORT=15672 +RABBITMQ_VIRTUAL_HOST=/ + +# rustfs (S3-compatible object storage) +RUSTFS_ACCESS_KEY=flexmodel +RUSTFS_SECRET_KEY=flexmodel123456 +RUSTFS_BUCKET=flexmodel +# S3 API port (host) - used by external tools; server reaches rustfs internally +RUSTFS_S3_API_PORT=9000 +# Web console port (host) +RUSTFS_CONSOLE_PORT=9001 + +# Project base domain (used for subdomain routing) +# e.g. preview.flexmodel.dev → {projectId}.preview.flexmodel.dev +FLEXMODEL_PROJECT_BASE_DOMAIN=flexmodel.wetech.tech \ No newline at end of file diff --git a/deploy/docker-compose/002-mysql-fullstack/.gitignore b/deploy/docker-compose/002-mysql-fullstack/.gitignore new file mode 100644 index 00000000..a349b226 --- /dev/null +++ b/deploy/docker-compose/002-mysql-fullstack/.gitignore @@ -0,0 +1,7 @@ +/mysql/data/ +/rustfs/data/ +/rabbitmq/data/ +/flexmodel/storage/ +/flexmodel/pages/ +/flexmodel/logs/ +/flexmodel/database/ diff --git a/deploy/docker-compose/002-mysql-fullstack/README.md b/deploy/docker-compose/002-mysql-fullstack/README.md new file mode 100644 index 00000000..214c4bb4 --- /dev/null +++ b/deploy/docker-compose/002-mysql-fullstack/README.md @@ -0,0 +1,66 @@ +# docker-compose + +Flexmodel 生产部署,包含以下服务: + +| 服务 | 镜像 | 端口 | +|-------------------------------|-------------------------------------------|--------------| +| `mysql` | `mysql:8.0` | 3306 (内部) | +| `rabbitmq` | `rabbitmq:3.13-management` | 5672 / 15672 | +| `rustfs` | `rustfs/rustfs:latest` | 9000 / 9001 | +| `flexmodel` | `cjbi/flexmodel-server:latest` | 8080 (内部) | +| `flexmodel-ui` | `cjbi/flexmodel-ui:latest` | 80 | +| `flexmodel-functions-runtime` | `cjbi/flexmodel-functions-runtime:latest` | 9999 (内部) | +| `nginx` | `nginx:1.25.3` | 80 (内部) | + +## 部署命令 + +* 拉取最新镜像 + +```shell +docker-compose pull +``` + +* 启动 + +```shell +docker-compose up +``` + +* 后台启动 + +```shell +docker-compose up -d +``` + +* 停止 + +```shell +docker-compose down +``` + +## RustFS 对象存储 + +本部署使用 [RustFS](https://rustfs.com)(S3 兼容的高性能对象存储)作为 `flexmodel-server` 的存储后端, 替代默认的本地磁盘存储。 +`flexmodel.storage.type` 被置为 `s3`,所有 Bucket / 文件操作经由 S3 API 写入 RustFS。 +| 服务 | 说明 | +|----------|---------------------------------------------------------------------| | `rustfs` | RustFS 主服务,S3 API +`9000`,Web 控制台 `9001`,使用宿主机目录持久化数据(./rustfs/data) | + +### 关键配置(`.env`) + +| 变量 | 说明 | 默认值 | +|-----------------------|-----------------------------|-------------------| +| `RUSTFS_ACCESS_KEY` | S3 Access Key | `flexmodel` | +| `RUSTFS_SECRET_KEY` | S3 Secret Key | `flexmodel123456` | +| `RUSTFS_BUCKET` | flexmodel-server 使用的桶名 | `flexmodel` | +| `RUSTFS_S3_API_PORT` | 对外暴露的 S3 API 端口 | `9000` | +| `RUSTFS_CONSOLE_PORT` | 对外暴露的 Web 控制台端口 | `9001` | + +### 注意事项 + +- 单盘单节点模式下已开启 `RUSTFS_UNSAFE_BYPASS_DISK_CHECK=true`;生产环境建议挂载多盘并关闭该选项以启用纠删码。 +- `flexmodel-server` 依赖 `rustfs` 健康检查通过后启动;启动时若桶不存在会自动创建(见下条)。 +- `flexmodel-server` 的 `S3Backend` 对所有 S3 兼容存储通用:启动校验时桶不存在则自动 `createBucket` + ,只读模式下不自动创建并报错提示运维预先建桶。因此本部署无需额外的初始化容器。 +- Web 控制台访问:`http://:${RUSTFS_CONSOLE_PORT}`,使用 `RUSTFS_ACCESS_KEY` / `RUSTFS_SECRET_KEY` 登录。 +- 切回本地存储:删除 `flexmodel-server` 环境变量中的 `FLEXMODEL_STORAGE_TYPE` 及 `QUARKUS_S3_*` 配置即可回退到本地存储。 diff --git a/deploy/docker-compose/002-mysql-fullstack/docker-compose.yml b/deploy/docker-compose/002-mysql-fullstack/docker-compose.yml new file mode 100644 index 00000000..7cca0220 --- /dev/null +++ b/deploy/docker-compose/002-mysql-fullstack/docker-compose.yml @@ -0,0 +1,177 @@ +services: + mysql: + container_name: ${CONTAINER_NAME_PREFIX}-flexmodel-mysql + image: mysql:8.0 + volumes: + - ./mysql/data:/var/lib/mysql + - ./mysql/init:/docker-entrypoint-initdb.d + - ./mysql/my.cnf:/etc/my.cnf + restart: always + environment: + - MYSQL_DATABASE=${MYSQL_DATABASE} + - MYSQL_ROOT_PASSWORD=${MYSQL_ROOT_PASSWORD} + - TZ=GMT+08 + command: > + bash -c " + chmod 644 /etc/my.cnf + && /entrypoint.sh mysqld + " + --default-authentication-plugin=mysql_native_password + --character-set-server=utf8mb4 + networks: + - flexmodel + healthcheck: + test: [ "CMD", "mysqladmin" ,"ping", "-h", "localhost" ] + timeout: 20s + retries: 10 + rabbitmq: + container_name: ${CONTAINER_NAME_PREFIX}-flexmodel-rabbitmq + image: rabbitmq:3.13-management + volumes: + - ./rabbitmq/data:/var/lib/rabbitmq + restart: always + environment: + - RABBITMQ_DEFAULT_USER=${RABBITMQ_DEFAULT_USER} + - RABBITMQ_DEFAULT_PASS=${RABBITMQ_DEFAULT_PASS} + - TZ=GMT+08 + ports: + - ${RABBITMQ_AMQP_PORT:-5672}:5672 + - ${RABBITMQ_MANAGEMENT_PORT:-15672}:15672 + networks: + - flexmodel + healthcheck: + test: [ "CMD", "rabbitmq-diagnostics", "-q", "ping" ] + timeout: 20s + retries: 10 + rustfs: + container_name: ${CONTAINER_NAME_PREFIX}-flexmodel-rustfs + image: rustfs/rustfs:latest + restart: always + volumes: + - ./rustfs/data:/data + environment: + - TZ=GMT+08 + - RUSTFS_VOLUMES=/data + - RUSTFS_ADDRESS=0.0.0.0:9000 + - RUSTFS_CONSOLE_ADDRESS=0.0.0.0:9001 + - RUSTFS_CONSOLE_ENABLE=true + - RUSTFS_CONSOLE_CORS_ALLOWED_ORIGINS=* + - RUSTFS_ACCESS_KEY=${RUSTFS_ACCESS_KEY} + - RUSTFS_SECRET_KEY=${RUSTFS_SECRET_KEY} + - RUSTFS_OBS_LOGGER_LEVEL=info + # single-disk single-node mode; bypass strict disk topology check + - RUSTFS_UNSAFE_BYPASS_DISK_CHECK=true + ports: + - ${RUSTFS_S3_API_PORT:-9000}:9000 + - ${RUSTFS_CONSOLE_PORT:-9001}:9001 + networks: + - flexmodel + healthcheck: + test: [ "CMD", "sh", "-ec", "curl -fsS http://127.0.0.1:9000/health" ] + interval: 30s + timeout: 10s + retries: 5 + start_period: 30s + flexmodel-ui: + container_name: ${CONTAINER_NAME_PREFIX}-flexmodel-ui + image: cjbi/flexmodel-ui:latest + restart: always + environment: + - TZ=GMT+08 + depends_on: + flexmodel-server: + condition: service_started + networks: + - flexmodel + nginx: + container_name: ${CONTAINER_NAME_PREFIX}-flexmodel-nginx + image: nginx:1.25.3 + volumes: + - ./nginx/nginx.conf:/etc/nginx/nginx.conf + - ./nginx/conf.d/default.conf.template:/etc/nginx/conf.d/default.conf.template + - ./flexmodel/pages:/data/pages + restart: always + ports: + - ${NGINX_HTTP_PORT:-80}:80 + - ${NGINX_HTTPS_PORT:-443}:443 + environment: + - TZ=GMT+08 + - FLEXMODEL_PROJECT_BASE_DOMAIN=${FLEXMODEL_PROJECT_BASE_DOMAIN:-localhost} + command: > + /bin/sh -c "export FLEXMODEL_PROJECT_BASE_DOMAIN_ESCAPED=$$(echo $$FLEXMODEL_PROJECT_BASE_DOMAIN | sed 's/\\./\\\\./g') && envsubst '$$FLEXMODEL_PROJECT_BASE_DOMAIN $$FLEXMODEL_PROJECT_BASE_DOMAIN_ESCAPED' < /etc/nginx/conf.d/default.conf.template > /etc/nginx/conf.d/default.conf && nginx -g 'daemon off;'" + depends_on: + flexmodel-server: + condition: service_started + networks: + - flexmodel + flexmodel-server: + container_name: ${CONTAINER_NAME_PREFIX}-flexmodel-server + image: ${FLEXMODEL_SERVER_IMAGE:-cjbi/flexmodel-server-native:latest} + restart: always + volumes: + - ./flexmodel/storage:/deployments/storage + - ./flexmodel/database:/deployments/database + - ./flexmodel/logs:/deployments/logs + - ./flexmodel/pages:/deployments/pages + env_file: .env + environment: + - TZ=GMT+08 + - FLEXMODEL_DATASOURCE_DB_KIND=mysql + # readTimeout: MySQL 容器重启后池中残留连接会读到已死对端,无超时会无限阻塞(曾导致事件循环冻结) + - FLEXMODEL_DATASOURCE_USERNAME=root + - FLEXMODEL_DATASOURCE_PASSWORD=${MYSQL_ROOT_PASSWORD} + - FLEXMODEL_PROJECT_URL_TEMPLATE=jdbc:mysql://mysql:3306/{{databaseName}}?connectTimeout=5000&readTimeout=15000 + - QUARKUS_REST_CLIENT_FUNCTION_RUNTIME_URL=http://flexmodel-functions-runtime:9999 + - QUARKUS_HTTP_LIMITS_MAX_BODY_SIZE=100M + - FLEXMODEL_PAGES_ROOT=/data/pages + - FLEXMODEL_PAGES_URL_TEMPLATE=${FLEXMODEL_PAGES_URL_TEMPLATE:-https://{{projectId}}.example.com} + - FLEXMODEL_EDGE_URL_TEMPLATE=${FLEXMODEL_EDGE_URL_TEMPLATE:-https://{{projectId}}.example.com/functions/{{name}}} + # RabbitMQ + - QUARKUS_RABBITMQ_HOST=rabbitmq + - QUARKUS_RABBITMQ_PORT=5672 + - QUARKUS_RABBITMQ_USERNAME=${RABBITMQ_DEFAULT_USER} + - QUARKUS_RABBITMQ_PASSWORD=${RABBITMQ_DEFAULT_PASS} + - QUARKUS_RABBITMQ_VIRTUAL_HOST=${RABBITMQ_VIRTUAL_HOST:-/} + - FLEXMODEL_EVENTS_RABBITMQ_ENABLED=${FLEXMODEL_RABBITMQ_ENABLED:-true} + + # RustFS (S3-compatible object storage) + - FLEXMODEL_STORAGE_TYPE=s3 + - FLEXMODEL_STORAGE_S3_BUCKET=${RUSTFS_BUCKET} + - FLEXMODEL_STORAGE_S3_ENDPOINT=http://rustfs:9000 + - FLEXMODEL_STORAGE_S3_REGION=us-east-1 + - FLEXMODEL_STORAGE_S3_PATH_STYLE=true + - QUARKUS_S3_AWS_REGION=us-east-1 + - QUARKUS_S3_AWS_CREDENTIALS_TYPE=static + - QUARKUS_S3_AWS_CREDENTIALS_STATIC_PROVIDER_ACCESS_KEY_ID=${RUSTFS_ACCESS_KEY} + - QUARKUS_S3_AWS_CREDENTIALS_STATIC_PROVIDER_SECRET_ACCESS_KEY=${RUSTFS_SECRET_KEY} + - QUARKUS_S3_ENDPOINT_OVERRIDE=http://rustfs:9000 + - QUARKUS_S3_PATH_STYLE_ACCESS=true + + networks: + - flexmodel + depends_on: + mysql: + condition: service_healthy + rabbitmq: + condition: service_healthy + rustfs: + condition: service_healthy + + flexmodel-functions-runtime: + container_name: ${CONTAINER_NAME_PREFIX}-flexmodel-functions-runtime + image: cjbi/flexmodel-functions-runtime:latest + restart: always + env_file: .env + environment: + - FLEXMODEL_JAVA_HOST=flexmodel-server + - FLEXMODEL_JAVA_PORT=8080 + - FLEXMODEL_PORT=9999 + - FLEXMODEL_HOST=0.0.0.0 + networks: + - flexmodel + depends_on: + - flexmodel-server +networks: + flexmodel: + name: ${CONTAINER_NAME_PREFIX}_flexmodel_default + driver: bridge diff --git a/deploy/docker-compose/002-mysql-fullstack/mysql/init/import.sql b/deploy/docker-compose/002-mysql-fullstack/mysql/init/import.sql new file mode 100644 index 00000000..e69de29b diff --git a/deploy/docker-compose/002-mysql-fullstack/mysql/my.cnf b/deploy/docker-compose/002-mysql-fullstack/mysql/my.cnf new file mode 100644 index 00000000..90e7e4ad --- /dev/null +++ b/deploy/docker-compose/002-mysql-fullstack/mysql/my.cnf @@ -0,0 +1,7 @@ +[client] +default-character-set=utf8mb4 +[mysql] +default-character-set=utf8mb4 +[mysqld] +character-set-server=utf8mb4 +init_connect='SET NAMES utf8mb4' diff --git a/deploy/docker-compose/002-mysql-fullstack/nginx/conf.d/default.conf.template b/deploy/docker-compose/002-mysql-fullstack/nginx/conf.d/default.conf.template new file mode 100644 index 00000000..eed82226 --- /dev/null +++ b/deploy/docker-compose/002-mysql-fullstack/nginx/conf.d/default.conf.template @@ -0,0 +1,105 @@ +# ============================================================ +# Pages alias subpath segment — avoids "//" when alias is empty. +# Evaluated lazily from $alias captured by the subdomain server_name. +# ============================================================ +map $alias $alias_path { + "" ""; + default "/$alias"; +} + +# ============================================================ +# Main domain → Admin API + UI +# /api/ → Java (admin API: /api/projects/... etc.) +# / → UI (flexmodel-ui) +# ============================================================ +server { + listen 80; + server_name localhost ${FLEXMODEL_PROJECT_BASE_DOMAIN}; + + client_max_body_size 1024m; + + # Admin API → Java + location /api/ { + proxy_pass http://flexmodel-server:8080; + proxy_set_header Host $host; + proxy_set_header X-Real-IP $remote_addr; + proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; + proxy_set_header REMOTE-HOST $remote_addr; + proxy_set_header Upgrade $http_upgrade; + proxy_set_header Connection "upgrade"; + proxy_http_version 1.1; + add_header Cache-Control no-cache; + } + + # Pages (path mode) — /pages/{projectId}/... → /data/pages/{projectId}/production/... + # 当 project-base-domain 为空时自动走 path 模式。 + location ~ ^/pages/(?[^/]+)(?.*)$ { + root /data/pages; + try_files /$pproj/production$prest /$pproj/production/index.html =404; + } + + # UI reverse proxy + location / { + proxy_pass http://flexmodel-ui:80; + proxy_set_header Host $host; + proxy_set_header X-Real-IP $remote_addr; + proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; + proxy_set_header Upgrade $http_upgrade; + proxy_set_header Connection "upgrade"; + proxy_http_version 1.1; + } + + error_page 500 502 503 504 /50x.html; + location = /50x.html { + root /usr/share/nginx/html; + } +} + +# ============================================================ +# Subdomain → Pages + Open API (data plane) +# {projectId}.{domain} → pages (production) +# {alias}.{projectId}.{domain} → pages (alias) +# /api/open/{projectId}/... → Java (open API) +# ============================================================ +server { + listen 80; + server_name ~^(?[^.]+)\.(?[^.]+)\.${FLEXMODEL_PROJECT_BASE_DOMAIN_ESCAPED}$ + ~^(?[^.]+)\.${FLEXMODEL_PROJECT_BASE_DOMAIN_ESCAPED}$; + + root /data/pages; + client_max_body_size 1024m; + + # Open API → Java (SDK sends /api/open/{projectId}/..., no rewrite needed) + location /api/ { + proxy_pass http://flexmodel-server:8080; + proxy_set_header Host $host; + proxy_set_header X-Real-IP $remote_addr; + proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; + proxy_set_header REMOTE-HOST $remote_addr; + proxy_set_header Upgrade $http_upgrade; + proxy_set_header Connection "upgrade"; + proxy_http_version 1.1; + add_header Cache-Control no-cache; + } + + # Pages base path (projectId + optional alias segment) + set $pages_base $projectId$alias_path; + + # Static assets — long cache (hash filenames) + location ~* \.(js|css|png|jpg|jpeg|gif|svg|woff2?|ttf|ico|mjs)$ { + try_files /$pages_base$uri /$projectId/production$uri =404; + expires 30d; + add_header Cache-Control "public, immutable"; + } + + # Root → serve index.html directly (avoids directory match → 403) + location = / { + try_files /$pages_base/index.html /$projectId/production/index.html =404; + } + + # Pages — SPA fallback (for non-root paths) + location / { + try_files /$pages_base$uri /$pages_base/index.html /$projectId/production$uri /$projectId/production/index.html =404; + } + +} diff --git a/deploy/docker-compose/002-mysql-fullstack/nginx/nginx.conf b/deploy/docker-compose/002-mysql-fullstack/nginx/nginx.conf new file mode 100644 index 00000000..7eabb059 --- /dev/null +++ b/deploy/docker-compose/002-mysql-fullstack/nginx/nginx.conf @@ -0,0 +1,35 @@ +user nginx; +worker_processes auto; + +error_log /var/log/nginx/error.log notice; +pid /var/run/nginx.pid; + + +events { + worker_connections 1024; +} + + +http { + include /etc/nginx/mime.types; + default_type application/octet-stream; + + log_format main '$remote_addr - $remote_user [$time_local] "$request" ' + '$status $body_bytes_sent "$http_referer" ' + '"$http_user_agent" "$http_x_forwarded_for"'; + + access_log /var/log/nginx/access.log main; + + sendfile on; + #tcp_nopush on; + + proxy_connect_timeout 600; + proxy_send_timeout 600; + proxy_read_timeout 600; + send_timeout 600; + keepalive_timeout 600; + + #gzip on; + + include /etc/nginx/conf.d/*.conf; +} diff --git a/feature_list.json b/feature_list.json index c782c4c0..7fff1fb6 100644 --- a/feature_list.json +++ b/feature_list.json @@ -81,6 +81,16 @@ "dependencies": ["feat-001", "feat-002"], "status": "in-progress", "evidence": "重构完成: 数据模型简化 (去除 f_function_version 表/branch字段/TriggerType.HTTP),Java 后端简化 (14 files — 删除4个版本/内部资源文件,重写Service/Resource/Invoker/DTO),Deno Functions Runtime 重写 (file:// 部署模式替代 lazy load),前端重写 (Supabase-style 多文件编辑器+模板选择+简化API)。Java 424 files + TypeScript 0 errors 编译通过。待: 端到端集成测试 + Deno 运行时验证。" + }, + { + "id": "feat-011", + "name": "Flow - 生命周期事件(本地 EventBus + 可选 RabbitMQ 桥接)", + "description": "在 flow 核心关键生命周期点发布 11 种强类型本地事件(FlowCreated/Updated/Deployed/Deleted、FlowInstanceStarted/Completed/Failed/Terminated、UserTaskSuspended/Committed/RollbackSuspended),经 Vert.x EventBus 广播;新增可选 RabbitMQ 桥接 Bean,默认关闭、未配置不连 broker,本地事件照常发布。核心与 RabbitMQ 解耦,本地事件可被多个内部订阅者复用。", + "dependencies": [ + "feat-003" + ], + "status": "done", + "evidence": "新增 dev.flexmodel.flow.event 包:FlowEvent 基类 + 11 个具体事件 + FlowEventTypes 常量;FlowEventPublisher 本地发布;FlowEventRabbitmqBridge(Instance 延迟注入,转发由 SmallRye 通道 mp.messaging.outgoing.events-out.enabled 单一控制,默认 false 注入 no-op emitter 不连 broker)。埋点 DefinitionProcessor/FlowExecutor/RuntimeProcessor.terminateProcess/UserTaskExecutor。application.properties 通道默认关闭+devservices=false。验证:mvn clean compile 通过;FlowEventPublisherTest(2)/FlowEventRabbitmqBridgeTest(2) 通过;DefinitionProcessorTest(4)/RuntimeProcessorTest(17) 回归通过。" } ] } diff --git a/flexmodel-sdks b/flexmodel-sdks index cb79787b..adcdf42c 160000 --- a/flexmodel-sdks +++ b/flexmodel-sdks @@ -1 +1 @@ -Subproject commit cb79787bd5b8a861bb18e8f3ed78e65d976e102e +Subproject commit adcdf42cfdf684c25a801d636699cc043b7b65c5 diff --git a/flexmodel-server/pom.xml b/flexmodel-server/pom.xml index e872c064..858b7fec 100644 --- a/flexmodel-server/pom.xml +++ b/flexmodel-server/pom.xml @@ -135,6 +135,20 @@ io.roastedroot quickjs4j-scripting-experimental + + + + io.quarkus + quarkus-messaging-rabbitmq + + + + org.testcontainers + rabbitmq + ${testcontainers.version} + test + + diff --git a/flexmodel-server/src/main/docker/Dockerfile.native b/flexmodel-server/src/main/docker/Dockerfile.native index ea38eb7e..76377bea 100644 --- a/flexmodel-server/src/main/docker/Dockerfile.native +++ b/flexmodel-server/src/main/docker/Dockerfile.native @@ -17,11 +17,11 @@ # To use UBI 8, switch to `quay.io/ubi8/ubi-minimal:8.10`. ### FROM registry.access.redhat.com/ubi9/ubi-minimal:9.7 -WORKDIR /work/ -RUN chown 1001 /work \ - && chmod "g+rwX" /work \ - && chown 1001:root /work -COPY --chown=1001:root --chmod=0755 target/*-runner /work/application +WORKDIR /deployments/ +RUN chown 1001 /deployments \ + && chmod "g+rwX" /deployments \ + && chown 1001:root /deployments +COPY --chown=1001:root --chmod=0755 target/*-runner /deployments/application EXPOSE 8080 USER 1001 diff --git a/flexmodel-server/src/main/docker/Dockerfile.native-micro b/flexmodel-server/src/main/docker/Dockerfile.native-micro index 63a5976e..28e17d8e 100644 --- a/flexmodel-server/src/main/docker/Dockerfile.native-micro +++ b/flexmodel-server/src/main/docker/Dockerfile.native-micro @@ -20,11 +20,11 @@ # To use UBI 8, switch to `quay.io/quarkus/quarkus-micro-image:2.0`. ### FROM quay.io/quarkus/ubi9-quarkus-micro-image:2.0 -WORKDIR /work/ -RUN chown 1001 /work \ - && chmod "g+rwX" /work \ - && chown 1001:root /work -COPY --chown=1001:root --chmod=0755 target/*-runner /work/application +WORKDIR /deployments/ +RUN chown 1001 /deployments \ + && chmod "g+rwX" /deployments \ + && chown 1001:root /deployments +COPY --chown=1001:root --chmod=0755 target/*-runner /deployments/application EXPOSE 8080 USER 1001 diff --git a/flexmodel-server/src/main/java/dev/flexmodel/FlexmodelApplication.java b/flexmodel-server/src/main/java/dev/flexmodel/FlexmodelApplication.java index 55857394..8697189e 100644 --- a/flexmodel-server/src/main/java/dev/flexmodel/FlexmodelApplication.java +++ b/flexmodel-server/src/main/java/dev/flexmodel/FlexmodelApplication.java @@ -38,8 +38,12 @@ public int run(String... args) throws Exception { private void printFlexmodelConfig() { log.info("---------- Flexmodel Configuration ----------"); + log.info(" version: {}", flexmodelConfig.version()); log.info(" project-url-template: {}", flexmodelConfig.projectUrlTemplate()); + log.info(" project-base-domain: {}", flexmodelConfig.projectBaseDomain().orElse("")); + log.info(" routing-mode: {}", flexmodelConfig.isSubdomainRouting() ? "subdomain" : "path"); log.info(" api-root-path: {}", flexmodelConfig.apiRootPath()); + log.info(" pages.root-path: {}", flexmodelConfig.pages().rootPath()); if (flexmodelConfig.datasources() != null && !flexmodelConfig.datasources().isEmpty()) { flexmodelConfig.datasources().forEach((name, ds) -> { log.info(" datasource[{}]:", name); @@ -51,6 +55,10 @@ private void printFlexmodelConfig() { } else { log.info(" (no datasources configured)"); } + log.info(" jwt.secret: {}", flexmodelConfig.jwt().secret().isBlank() ? "" : "***"); + log.info(" jwt.access-token-lifetime: {}", flexmodelConfig.jwt().accessTokenLifetime()); + log.info(" jwt.refresh-token-lifetime: {}", flexmodelConfig.jwt().refreshTokenLifetime()); + log.info(" events.rabbitmq.enabled: {}", flexmodelConfig.events().rabbitmq().enabled()); } } diff --git a/flexmodel-server/src/main/java/dev/flexmodel/common/FlexmodelConfig.java b/flexmodel-server/src/main/java/dev/flexmodel/common/FlexmodelConfig.java index e8e19d15..1d9eaf78 100644 --- a/flexmodel-server/src/main/java/dev/flexmodel/common/FlexmodelConfig.java +++ b/flexmodel-server/src/main/java/dev/flexmodel/common/FlexmodelConfig.java @@ -92,6 +92,25 @@ interface PagesConfig { String rootPath(); } + /** + * 事件转发 RabbitMQ 桥接的业务开关。默认 false:不连 broker,桥接静默跳过。 + * 该开关经 application.properties 占位符映射到 SmallRye 通道 mp.messaging.outgoing.events-out.enabled。 + */ + @WithName("events") + EventsConfig events(); + + interface EventsConfig { + + @WithName("rabbitmq") + RabbitmqConfig rabbitmq(); + + interface RabbitmqConfig { + + @WithDefault("false") + boolean enabled(); + } + } + interface DatasourceConfig { @WithName("db-kind") diff --git a/flexmodel-server/src/main/java/dev/flexmodel/common/config/EngineConfig.java b/flexmodel-server/src/main/java/dev/flexmodel/common/config/EngineConfig.java index 0d4d0dcc..895203c2 100644 --- a/flexmodel-server/src/main/java/dev/flexmodel/common/config/EngineConfig.java +++ b/flexmodel-server/src/main/java/dev/flexmodel/common/config/EngineConfig.java @@ -8,6 +8,7 @@ import dev.flexmodel.project.BranchRepository; import dev.flexmodel.project.ProjectService; import dev.flexmodel.realtime.RealtimeEventListener; +import dev.flexmodel.realtime.RealtimeRabbitmqListener; import dev.flexmodel.scheduling.TriggerDataChangedEventListener; import dev.flexmodel.session.SessionFactory; import dev.flexmodel.sql.JdbcSchemaProvider; @@ -51,7 +52,8 @@ public void installDatasource(@Observes StartupEvent startupEvent, public SessionFactory sessionFactory(FlexmodelConfig flexmodelConfig, TriggerDataChangedEventListener triggerDataChangedEventListener, AuditDataEventListener auditDataEventListener, - RealtimeEventListener realtimeEventListener) { + RealtimeEventListener realtimeEventListener, + RealtimeRabbitmqListener realtimeRabbitmqListener) { FlexmodelConfig.DatasourceConfig datasourceConfig = flexmodelConfig.datasources().get(SYSTEM_DS_KEY); AgroalDataSource defaultDs = AgroalDataSourceFactory.createDataSource( datasourceConfig.url(), @@ -80,6 +82,7 @@ public SessionFactory sessionFactory(FlexmodelConfig flexmodelConfig, sf.getEventPublisher().addListener(triggerDataChangedEventListener); sf.getEventPublisher().addListener(auditDataEventListener); sf.getEventPublisher().addListener(realtimeEventListener); + sf.getEventPublisher().addListener(realtimeRabbitmqListener); return sf; } diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowEvent.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowEvent.java new file mode 100644 index 00000000..49bb14cf --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowEvent.java @@ -0,0 +1,66 @@ +package dev.flexmodel.flow.event; + +import java.io.Serializable; + +/** + * Flow 生命周期事件基类。 + *

+ * 公共字段 {@code projectId}、{@code caller}、{@code timestamp}(构造时设为当前毫秒)。 + * 抽象方法 {@link #routingKey()} 作 Vert.x EventBus 地址(静态常量,供 {@code @ConsumeEvent} 匹配)。 + * {@link #rabbitmqRoutingKey()} 在 routing key 中插入 projectId,作 RabbitMQ 转发的 routing key, + * 格式为 {@code flow..},便于消费端按项目订阅。 + * 本地事件始终经 {@code FlowEventPublisher} 发布到 EventBus;RabbitMQ 转发可选。 + * + * @author cjbi + */ +public abstract class FlowEvent implements Serializable { + + private static final String FLOW_PREFIX = "flow."; + + private final String projectId; + private final String caller; + private final long timestamp; + + protected FlowEvent(String projectId, String caller) { + this.projectId = projectId; + this.caller = caller; + this.timestamp = System.currentTimeMillis(); + } + + /** + * EventBus 地址(静态常量),供 {@code @ConsumeEvent} 匹配。 + *

+ * 形如 {@code flow.instance.started},不含 projectId。 + * + * @return 不可为空的 EventBus 地址 + */ + public abstract String routingKey(); + + /** + * RabbitMQ routing key,在 {@link #routingKey()} 中插入 projectId。 + *

+ * 形如 {@code flow..instance.started},便于消费端按项目通配符订阅。 + * 若 routingKey 不以 {@code flow.} 开头则原样返回(防御性兜底)。 + * + * @return 带 projectId 的 RabbitMQ routing key + */ + public String rabbitmqRoutingKey() { + String key = routingKey(); + if (key.startsWith(FLOW_PREFIX)) { + return FLOW_PREFIX + projectId + "." + key.substring(FLOW_PREFIX.length()); + } + return key; + } + + public String getProjectId() { + return projectId; + } + + public String getCaller() { + return caller; + } + + public long getTimestamp() { + return timestamp; + } +} \ No newline at end of file diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowEventPublisher.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowEventPublisher.java new file mode 100644 index 00000000..a935b92b --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowEventPublisher.java @@ -0,0 +1,40 @@ +package dev.flexmodel.flow.event; + +import io.vertx.mutiny.core.eventbus.EventBus; +import jakarta.enterprise.context.ApplicationScoped; +import jakarta.inject.Inject; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Flow 生命周期事件本地发布器。 + *

+ * 始终将事件发布到 Vert.x EventBus(本地事件为机制本身),不与 RabbitMQ 配置耦合。 + * 全程 try/catch 仅 {@code LOGGER.warn},不抛出、不阻塞流程。 + * + * @author cjbi + */ +@ApplicationScoped +public class FlowEventPublisher { + + private static final Logger LOGGER = LoggerFactory.getLogger(FlowEventPublisher.class); + + @Inject + EventBus eventBus; + + /** + * 发布 Flow 生命周期事件。失败仅告警,不影响流程主链路。 + * + * @param event 不可为空的事件 + */ + public void publish(FlowEvent event) { + if (event == null) { + return; + } + try { + eventBus.publish(event.routingKey(), event); + } catch (Exception e) { + LOGGER.warn("publish flow event failed.||routingKey={}||event={}", event.routingKey(), event, e); + } + } +} diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowEventRabbitmqBridge.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowEventRabbitmqBridge.java new file mode 100644 index 00000000..c7bf42a3 --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowEventRabbitmqBridge.java @@ -0,0 +1,103 @@ +package dev.flexmodel.flow.event; + +import io.smallrye.reactive.messaging.MutinyEmitter; +import io.smallrye.reactive.messaging.rabbitmq.OutgoingRabbitMQMetadata; +import io.quarkus.vertx.ConsumeEvent; +import jakarta.enterprise.context.ApplicationScoped; +import jakarta.enterprise.inject.Instance; +import jakarta.inject.Inject; +import org.eclipse.microprofile.config.inject.ConfigProperty; +import org.eclipse.microprofile.reactive.messaging.Channel; +import org.eclipse.microprofile.reactive.messaging.Message; +import org.eclipse.microprofile.reactive.messaging.Metadata; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Flow 生命周期事件 RabbitMQ 桥接。 + *

+ * 订阅全部内部 EventBus 事件并转发到单个 topic 交换机(routing key 复用事件地址)。 + * 默认不连 broker({@code flexmodel.events.rabbitmq.enabled=false}): + * 关闭时 SmallRye 注入 no-op emitter,{@link #forward(FlowEvent)} 发送即丢弃、不连 broker; + * 本地事件照常经 {@link FlowEventPublisher} 发布。启用需置 {@code flexmodel.events.rabbitmq.enabled=true}。 + *

+ * 每个事件一个 {@code @ConsumeEvent} 方法、消费具体类型,沿用代码库已验证的同类型发布/消费模式, + * 避免多态反序列化不确定性。 + *

+ * 使用 {@link Instance} 延迟注入 {@link MutinyEmitter}:通道未启用时 SmallRye 提供 no-op emitter, + * bean 仍可正常创建,不影响本地事件广播。 + * + * @author cjbi + */ +@ApplicationScoped +public class FlowEventRabbitmqBridge { + + private static final Logger LOGGER = LoggerFactory.getLogger(FlowEventRabbitmqBridge.class); + + @Inject + @Channel("events-out") + Instance> flowEventEmitterInstance; + + @ConsumeEvent(value = FlowEventTypes.FLOW_INSTANCE_STARTED, blocking = false) + public void onFlowInstanceStarted(FlowInstanceStartedEvent event) { + forward(event); + } + + @ConsumeEvent(value = FlowEventTypes.FLOW_INSTANCE_COMPLETED, blocking = false) + public void onFlowInstanceCompleted(FlowInstanceCompletedEvent event) { + forward(event); + } + + @ConsumeEvent(value = FlowEventTypes.FLOW_INSTANCE_FAILED, blocking = false) + public void onFlowInstanceFailed(FlowInstanceFailedEvent event) { + forward(event); + } + + @ConsumeEvent(value = FlowEventTypes.FLOW_INSTANCE_TERMINATED, blocking = false) + public void onFlowInstanceTerminated(FlowInstanceTerminatedEvent event) { + forward(event); + } + + @ConsumeEvent(value = FlowEventTypes.USER_TASK_SUSPENDED, blocking = false) + public void onUserTaskSuspended(UserTaskSuspendedEvent event) { + forward(event); + } + + @ConsumeEvent(value = FlowEventTypes.USER_TASK_COMMITTED, blocking = false) + public void onUserTaskCommitted(UserTaskCommittedEvent event) { + forward(event); + } + + @ConsumeEvent(value = FlowEventTypes.USER_TASK_ROLLBACK_SUSPENDED, blocking = false) + public void onUserTaskRollbackSuspended(UserTaskRollbackSuspendedEvent event) { + forward(event); + } + + /** + * 尽力而为转发:以事件 routing key 作 RabbitMQ routing key,失败仅告警、不重试、不阻塞。 + * 通道未启用时 SmallRye 提供 no-op emitter,发送即丢弃、不连 broker。 + */ + private void forward(FlowEvent event) { + if (flowEventEmitterInstance.isUnsatisfied()) { + LOGGER.debug("events-out channel not resolvable, skip forwarding.||routingKey={}", event.rabbitmqRoutingKey()); + return; + } + try { + MutinyEmitter emitter = flowEventEmitterInstance.get(); + OutgoingRabbitMQMetadata metadata = new OutgoingRabbitMQMetadata.Builder() + .withRoutingKey(event.rabbitmqRoutingKey()) + // 持久化投递(delivery_mode=2):使消息在 broker 重启后仍可恢复,配合 durable 交换机与 durable 队列生效 + .withDeliveryMode(2) + .build(); + emitter.sendMessage(Message.of(event, Metadata.of(metadata))) + .onFailure() + .invoke(e -> LOGGER.warn("forward flow event to rabbitmq failed.||routingKey={}||payloadType={}", + event.rabbitmqRoutingKey(), event.getClass().getSimpleName(), e)) + .subscribe() + .asCompletionStage(); + } catch (Exception e) { + // 仅记录 routing key 与异常,不打印事件载荷(避免泄露流程变量) + LOGGER.warn("forward flow event to rabbitmq failed (sync).||routingKey={}", event.rabbitmqRoutingKey(), e); + } + } +} diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowEventTypes.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowEventTypes.java new file mode 100644 index 00000000..5f62652b --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowEventTypes.java @@ -0,0 +1,25 @@ +package dev.flexmodel.flow.event; + +/** + * Flow 生命周期事件路由 key 常量。 + *

+ * 供 {@code @ConsumeEvent} 注解与 RabbitMQ 桥接消费者引用,避免散落字符串字面量。 + * + * @author cjbi + */ +public final class FlowEventTypes { + + private FlowEventTypes() { + } + + // 实例层 + public static final String FLOW_INSTANCE_STARTED = "flow.instance.started"; + public static final String FLOW_INSTANCE_COMPLETED = "flow.instance.completed"; + public static final String FLOW_INSTANCE_FAILED = "flow.instance.failed"; + public static final String FLOW_INSTANCE_TERMINATED = "flow.instance.terminated"; + + // 用户任务层 + public static final String USER_TASK_SUSPENDED = "flow.usertask.suspended"; + public static final String USER_TASK_COMMITTED = "flow.usertask.committed"; + public static final String USER_TASK_ROLLBACK_SUSPENDED = "flow.usertask.rollback.suspended"; +} diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowInstanceCompletedEvent.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowInstanceCompletedEvent.java new file mode 100644 index 00000000..47a36976 --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowInstanceCompletedEvent.java @@ -0,0 +1,35 @@ +package dev.flexmodel.flow.event; + +import lombok.Getter; +import lombok.ToString; + +import java.util.Map; + +/** + * 流程实例完成事件。 + * + * @author cjbi + */ +@Getter +@ToString(callSuper = true) +public class FlowInstanceCompletedEvent extends FlowEvent { + + public static final String ROUTING_KEY = FlowEventTypes.FLOW_INSTANCE_COMPLETED; + + private final String flowDeployId; + private final String flowInstanceId; + private final Map variables; + + public FlowInstanceCompletedEvent(String projectId, String caller, String flowDeployId, + String flowInstanceId, Map variables) { + super(projectId, caller); + this.flowDeployId = flowDeployId; + this.flowInstanceId = flowInstanceId; + this.variables = variables == null ? null : new java.util.HashMap<>(variables); + } + + @Override + public String routingKey() { + return ROUTING_KEY; + } +} diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowInstanceFailedEvent.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowInstanceFailedEvent.java new file mode 100644 index 00000000..c4c04323 --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowInstanceFailedEvent.java @@ -0,0 +1,33 @@ +package dev.flexmodel.flow.event; + +import lombok.Getter; +import lombok.ToString; + +/** + * 流程实例失败事件。 + * + * @author cjbi + */ +@Getter +@ToString(callSuper = true) +public class FlowInstanceFailedEvent extends FlowEvent { + + public static final String ROUTING_KEY = FlowEventTypes.FLOW_INSTANCE_FAILED; + + private final String flowDeployId; + private final String flowInstanceId; + private final String error; + + public FlowInstanceFailedEvent(String projectId, String caller, String flowDeployId, + String flowInstanceId, String error) { + super(projectId, caller); + this.flowDeployId = flowDeployId; + this.flowInstanceId = flowInstanceId; + this.error = error; + } + + @Override + public String routingKey() { + return ROUTING_KEY; + } +} diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowInstanceStartedEvent.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowInstanceStartedEvent.java new file mode 100644 index 00000000..fd6f8a77 --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowInstanceStartedEvent.java @@ -0,0 +1,35 @@ +package dev.flexmodel.flow.event; + +import lombok.Getter; +import lombok.ToString; + +import java.util.Map; + +/** + * 流程实例启动事件。 + * + * @author cjbi + */ +@Getter +@ToString(callSuper = true) +public class FlowInstanceStartedEvent extends FlowEvent { + + public static final String ROUTING_KEY = FlowEventTypes.FLOW_INSTANCE_STARTED; + + private final String flowDeployId; + private final String flowInstanceId; + private final Map variables; + + public FlowInstanceStartedEvent(String projectId, String caller, String flowDeployId, + String flowInstanceId, Map variables) { + super(projectId, caller); + this.flowDeployId = flowDeployId; + this.flowInstanceId = flowInstanceId; + this.variables = variables == null ? null : new java.util.HashMap<>(variables); + } + + @Override + public String routingKey() { + return ROUTING_KEY; + } +} diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowInstanceTerminatedEvent.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowInstanceTerminatedEvent.java new file mode 100644 index 00000000..8068ac9d --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/FlowInstanceTerminatedEvent.java @@ -0,0 +1,28 @@ +package dev.flexmodel.flow.event; + +import lombok.Getter; +import lombok.ToString; + +/** + * 流程实例终止事件。 + * + * @author cjbi + */ +@Getter +@ToString(callSuper = true) +public class FlowInstanceTerminatedEvent extends FlowEvent { + + public static final String ROUTING_KEY = FlowEventTypes.FLOW_INSTANCE_TERMINATED; + + private final String flowInstanceId; + + public FlowInstanceTerminatedEvent(String projectId, String caller, String flowInstanceId) { + super(projectId, caller); + this.flowInstanceId = flowInstanceId; + } + + @Override + public String routingKey() { + return ROUTING_KEY; + } +} diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/StartFlowEvent.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/StartFlowEvent.java deleted file mode 100644 index f12c732b..00000000 --- a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/StartFlowEvent.java +++ /dev/null @@ -1,17 +0,0 @@ -package dev.flexmodel.flow.event; - -import lombok.Getter; -import lombok.Setter; - -import java.util.Map; - -/** - * @author cjbi - */ -@Getter -@Setter -public class StartFlowEvent { - private String flowDeployId; - private String flowModuleId; - private Map variables; -} diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/UserTaskCommittedEvent.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/UserTaskCommittedEvent.java new file mode 100644 index 00000000..c4cbb2d2 --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/UserTaskCommittedEvent.java @@ -0,0 +1,39 @@ +package dev.flexmodel.flow.event; + +import lombok.Getter; +import lombok.ToString; + +import java.util.Map; + +/** + * 用户任务提交完成事件。 + * + * @author cjbi + */ +@Getter +@ToString(callSuper = true) +public class UserTaskCommittedEvent extends FlowEvent { + + public static final String ROUTING_KEY = FlowEventTypes.USER_TASK_COMMITTED; + + private final String flowDeployId; + private final String flowInstanceId; + private final String nodeInstanceId; + private final String nodeKey; + private final Map nodeAttributes; + + public UserTaskCommittedEvent(String projectId, String caller, String flowDeployId, String flowInstanceId, + String nodeInstanceId, String nodeKey, Map nodeAttributes) { + super(projectId, caller); + this.flowDeployId = flowDeployId; + this.flowInstanceId = flowInstanceId; + this.nodeInstanceId = nodeInstanceId; + this.nodeKey = nodeKey; + this.nodeAttributes = nodeAttributes == null ? null : new java.util.HashMap<>(nodeAttributes); + } + + @Override + public String routingKey() { + return ROUTING_KEY; + } +} diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/UserTaskRollbackSuspendedEvent.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/UserTaskRollbackSuspendedEvent.java new file mode 100644 index 00000000..8e916c99 --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/UserTaskRollbackSuspendedEvent.java @@ -0,0 +1,39 @@ +package dev.flexmodel.flow.event; + +import lombok.Getter; +import lombok.ToString; + +import java.util.Map; + +/** + * 用户任务回滚挂起事件。 + * + * @author cjbi + */ +@Getter +@ToString(callSuper = true) +public class UserTaskRollbackSuspendedEvent extends FlowEvent { + + public static final String ROUTING_KEY = FlowEventTypes.USER_TASK_ROLLBACK_SUSPENDED; + + private final String flowDeployId; + private final String flowInstanceId; + private final String nodeInstanceId; + private final String nodeKey; + private final Map nodeAttributes; + + public UserTaskRollbackSuspendedEvent(String projectId, String caller, String flowDeployId, String flowInstanceId, + String nodeInstanceId, String nodeKey, Map nodeAttributes) { + super(projectId, caller); + this.flowDeployId = flowDeployId; + this.flowInstanceId = flowInstanceId; + this.nodeInstanceId = nodeInstanceId; + this.nodeKey = nodeKey; + this.nodeAttributes = nodeAttributes == null ? null : new java.util.HashMap<>(nodeAttributes); + } + + @Override + public String routingKey() { + return ROUTING_KEY; + } +} diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/event/UserTaskSuspendedEvent.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/UserTaskSuspendedEvent.java new file mode 100644 index 00000000..f61280ab --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/event/UserTaskSuspendedEvent.java @@ -0,0 +1,42 @@ +package dev.flexmodel.flow.event; + +import lombok.Getter; +import lombok.ToString; + +import java.util.Map; + +/** + * 用户任务挂起事件。 + * + * @author cjbi + */ +@Getter +@ToString(callSuper = true) +public class UserTaskSuspendedEvent extends FlowEvent { + + public static final String ROUTING_KEY = FlowEventTypes.USER_TASK_SUSPENDED; + + private final String flowDeployId; + private final String flowInstanceId; + private final String nodeInstanceId; + private final String nodeKey; + private final Map variables; + private final Map nodeAttributes; + + public UserTaskSuspendedEvent(String projectId, String caller, String flowDeployId, String flowInstanceId, + String nodeInstanceId, String nodeKey, Map variables, + Map nodeAttributes) { + super(projectId, caller); + this.flowDeployId = flowDeployId; + this.flowInstanceId = flowInstanceId; + this.nodeInstanceId = nodeInstanceId; + this.nodeKey = nodeKey; + this.variables = variables == null ? null : new java.util.HashMap<>(variables); + this.nodeAttributes = nodeAttributes == null ? null : new java.util.HashMap<>(nodeAttributes); + } + + @Override + public String routingKey() { + return ROUTING_KEY; + } +} diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/executor/FlowExecutor.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/executor/FlowExecutor.java index 36aad99a..e7867d3f 100644 --- a/flexmodel-server/src/main/java/dev/flexmodel/flow/executor/FlowExecutor.java +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/executor/FlowExecutor.java @@ -15,6 +15,10 @@ import dev.flexmodel.flow.exception.ProcessException; import dev.flexmodel.flow.exception.ReentrantException; import dev.flexmodel.flow.repository.FlowInstanceRepository; +import dev.flexmodel.flow.event.FlowEventPublisher; +import dev.flexmodel.flow.event.FlowInstanceCompletedEvent; +import dev.flexmodel.flow.event.FlowInstanceFailedEvent; +import dev.flexmodel.flow.event.FlowInstanceStartedEvent; import dev.flexmodel.flow.common.util.FlowModelUtil; import dev.flexmodel.flow.common.util.InstanceDataUtil; import dev.flexmodel.common.utils.CollectionUtils; @@ -38,11 +42,15 @@ public class FlowExecutor extends RuntimeExecutor { @Inject Instance executorFactoryInstance; + @Inject + FlowEventPublisher flowEventPublisher; + ////////////////////////////////////////execute//////////////////////////////////////// @Override public void execute(RuntimeContext runtimeContext) throws ProcessException { int processStatus = ProcessStatus.SUCCESS; + String errorMessage = null; try { preExecute(runtimeContext); doExecute(runtimeContext); @@ -50,10 +58,15 @@ public void execute(RuntimeContext runtimeContext) throws ProcessException { if (!ErrorEnum.isSuccess(pe.getErrNo())) { processStatus = ProcessStatus.FAILED; } + errorMessage = pe.getMessage(); throw pe; } finally { runtimeContext.setProcessStatus(processStatus); postExecute(runtimeContext); + if (processStatus == ProcessStatus.FAILED) { + flowEventPublisher.publish(new FlowInstanceFailedEvent(runtimeContext.getProjectId(), runtimeContext.getCaller(), + runtimeContext.getFlowDeployId(), runtimeContext.getFlowInstanceId(), errorMessage)); + } } } @@ -74,6 +87,9 @@ private void preExecute(RuntimeContext runtimeContext) throws ProcessException { //3.update runtimeContext fillExecuteContext(runtimeContext, flowInstancePO.getFlowInstanceId(), instanceDataId); + + flowEventPublisher.publish(new FlowInstanceStartedEvent(runtimeContext.getProjectId(), runtimeContext.getCaller(), + runtimeContext.getFlowDeployId(), runtimeContext.getFlowInstanceId(), runtimeContext.getInstanceDataMap())); } private FlowInstance saveFlowInstance(RuntimeContext runtimeContext) throws ProcessException { @@ -190,6 +206,9 @@ private void postExecute(RuntimeContext runtimeContext) throws ProcessException runtimeContext.setFlowInstanceStatus(FlowInstanceStatus.COMPLETED); } LOGGER.info("postExecute: flowInstance process completely.||flowInstanceId={}", runtimeContext.getFlowInstanceId()); + + flowEventPublisher.publish(new FlowInstanceCompletedEvent(runtimeContext.getProjectId(), runtimeContext.getCaller(), + runtimeContext.getFlowDeployId(), runtimeContext.getFlowInstanceId(), runtimeContext.getInstanceDataMap())); } } @@ -337,6 +356,9 @@ private void postCommit(RuntimeContext runtimeContext) throws ProcessException { } LOGGER.info("postCommit: flowInstance process completely.||flowInstanceId={}", runtimeContext.getFlowInstanceId()); + + flowEventPublisher.publish(new FlowInstanceCompletedEvent(runtimeContext.getProjectId(), runtimeContext.getCaller(), + runtimeContext.getFlowDeployId(), runtimeContext.getFlowInstanceId(), runtimeContext.getInstanceDataMap())); } } diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/executor/UserTaskExecutor.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/executor/UserTaskExecutor.java index d84b5ebf..e4e3c80c 100644 --- a/flexmodel-server/src/main/java/dev/flexmodel/flow/executor/UserTaskExecutor.java +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/executor/UserTaskExecutor.java @@ -1,5 +1,6 @@ package dev.flexmodel.flow.executor; +import jakarta.inject.Inject; import jakarta.inject.Singleton; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -12,6 +13,10 @@ import dev.flexmodel.flow.common.NodeInstanceStatus; import dev.flexmodel.flow.common.RuntimeContext; import dev.flexmodel.flow.common.util.FlowModelUtil; +import dev.flexmodel.flow.event.FlowEventPublisher; +import dev.flexmodel.flow.event.UserTaskCommittedEvent; +import dev.flexmodel.flow.event.UserTaskRollbackSuspendedEvent; +import dev.flexmodel.flow.event.UserTaskSuspendedEvent; import dev.flexmodel.JsonUtils; import java.text.MessageFormat; @@ -22,6 +27,9 @@ public class UserTaskExecutor extends ElementExecutor { private static final Logger LOGGER = LoggerFactory.getLogger(UserTaskExecutor.class); + @Inject + FlowEventPublisher flowEventPublisher; + @Override protected void doExecute(RuntimeContext runtimeContext) throws ProcessException { NodeInstanceBO currentNodeInstance = runtimeContext.getCurrentNodeInstance(); @@ -39,6 +47,10 @@ protected void doExecute(RuntimeContext runtimeContext) throws ProcessException String nodeName = FlowModelUtil.getElementName(flowElement); LOGGER.info("doExecute: userTask to commit.||flowInstanceId={}||nodeInstanceId={}||nodeKey={}||nodeName={}", runtimeContext.getFlowInstanceId(), currentNodeInstance.getNodeInstanceId(), flowElement.getKey(), nodeName); + flowEventPublisher.publish(new UserTaskSuspendedEvent(runtimeContext.getProjectId(), runtimeContext.getCaller(), + runtimeContext.getFlowDeployId(), runtimeContext.getFlowInstanceId(), + currentNodeInstance.getNodeInstanceId(), flowElement.getKey(), runtimeContext.getInstanceDataMap(), + flowElement.getProperties())); throw new SuspendException(ErrorEnum.COMMIT_SUSPEND, MessageFormat.format(Constants.NODE_INSTANCE_FORMAT, flowElement.getKey(), nodeName, currentNodeInstance.getNodeInstanceId())); } @@ -95,6 +107,11 @@ protected void postCommit(RuntimeContext runtimeContext) { if (currentNodeInstance.getStatus() != NodeInstanceStatus.COMPLETED) { currentNodeInstance.setStatus(NodeInstanceStatus.COMPLETED); runtimeContext.getNodeInstanceList().add(currentNodeInstance); + FlowElement committedElement = runtimeContext.getCurrentNodeModel(); + flowEventPublisher.publish(new UserTaskCommittedEvent(runtimeContext.getProjectId(), runtimeContext.getCaller(), + runtimeContext.getFlowDeployId(), runtimeContext.getFlowInstanceId(), + currentNodeInstance.getNodeInstanceId(), currentNodeInstance.getNodeKey(), + committedElement == null ? null : committedElement.getProperties())); } } @@ -123,6 +140,11 @@ protected void doRollback(RuntimeContext runtimeContext) throws ProcessException newNodeInstanceBO.setProperties(currentNodeInstance.getProperties()); runtimeContext.setCurrentNodeInstance(newNodeInstanceBO); runtimeContext.getNodeInstanceList().add(newNodeInstanceBO); + FlowElement rollbackElement = FlowModelUtil.getFlowElement(runtimeContext.getFlowElementMap(), newNodeInstanceBO.getNodeKey()); + flowEventPublisher.publish(new UserTaskRollbackSuspendedEvent(runtimeContext.getProjectId(), runtimeContext.getCaller(), + runtimeContext.getFlowDeployId(), runtimeContext.getFlowInstanceId(), + newNodeInstanceBO.getNodeInstanceId(), newNodeInstanceBO.getNodeKey(), + rollbackElement == null ? null : rollbackElement.getProperties())); throw new SuspendException(ErrorEnum.ROLLBACK_SUSPEND, MessageFormat.format(Constants.NODE_INSTANCE_FORMAT, newNodeInstanceBO.getNodeKey(), FlowModelUtil.getFlowElement(runtimeContext.getFlowElementMap(), newNodeInstanceBO.getNodeKey()), diff --git a/flexmodel-server/src/main/java/dev/flexmodel/flow/processor/RuntimeProcessor.java b/flexmodel-server/src/main/java/dev/flexmodel/flow/processor/RuntimeProcessor.java index 678114ac..0ccb04f0 100644 --- a/flexmodel-server/src/main/java/dev/flexmodel/flow/processor/RuntimeProcessor.java +++ b/flexmodel-server/src/main/java/dev/flexmodel/flow/processor/RuntimeProcessor.java @@ -24,6 +24,8 @@ import dev.flexmodel.flow.exception.ReentrantException; import dev.flexmodel.flow.exception.TurboException; import dev.flexmodel.flow.executor.FlowExecutor; +import dev.flexmodel.flow.event.FlowEventPublisher; +import dev.flexmodel.flow.event.FlowInstanceTerminatedEvent; import dev.flexmodel.flow.repository.FlowDeploymentRepository; import dev.flexmodel.flow.repository.FlowInstanceMappingRepository; import dev.flexmodel.flow.repository.FlowInstanceRepository; @@ -69,6 +71,9 @@ public class RuntimeProcessor { @Inject NodeInstanceService nodeInstanceService; + @Inject + FlowEventPublisher flowEventPublisher; + ////////////////////////////////////////startProcess//////////////////////////////////////// public StartProcessResult startProcess(StartProcessParam startProcessParam) { @@ -253,6 +258,7 @@ public TerminateResult terminateProcess(String projectId, String flowInstanceId, } else { processInstanceRepository.updateStatus(projectId, flowInstancePO, FlowInstanceStatus.TERMINATED); flowInstanceStatus = FlowInstanceStatus.TERMINATED; + flowEventPublisher.publish(new FlowInstanceTerminatedEvent(projectId, flowInstancePO.getCaller(), flowInstanceId)); } if (effectiveForSubFlowInstance) { terminateSubFlowInstance(projectId, flowInstanceId); diff --git a/flexmodel-server/src/main/java/dev/flexmodel/realtime/DataChangeEvent.java b/flexmodel-server/src/main/java/dev/flexmodel/realtime/DataChangeEvent.java new file mode 100644 index 00000000..9a835628 --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/realtime/DataChangeEvent.java @@ -0,0 +1,81 @@ +package dev.flexmodel.realtime; + +import java.io.Serializable; +import java.util.Map; + +/** + * 数据变更事件载荷,由 {@link RealtimeRabbitmqListener} 从引擎 {@code ChangedEvent} 转换而来, + * 经 SmallRye RabbitMQ 通道 {@code data-events-out} 投递到 {@code flexmodel.events} topic 交换机。 + *

+ * routing key 形如 {@code data...},便于消费端按项目、模型或操作类型订阅过滤。 + * + * @author cjbi + */ +public class DataChangeEvent implements Serializable { + + private final String routingKey; + private final String projectId; + private final String event; + private final String model; + private final String schema; + private final Object recordId; + private final long timestamp; + private final int affectedRows; + private final Map data; + private final Map oldData; + + public DataChangeEvent(String routingKey, String projectId, String event, String model, String schema, + Object recordId, long timestamp, int affectedRows, + Map data, Map oldData) { + this.routingKey = routingKey; + this.projectId = projectId; + this.event = event; + this.model = model; + this.schema = schema; + this.recordId = recordId; + this.timestamp = timestamp; + this.affectedRows = affectedRows; + this.data = data; + this.oldData = oldData; + } + + public String getRoutingKey() { + return routingKey; + } + + public String getProjectId() { + return projectId; + } + + public String getEvent() { + return event; + } + + public String getModel() { + return model; + } + + public String getSchema() { + return schema; + } + + public Object getRecordId() { + return recordId; + } + + public long getTimestamp() { + return timestamp; + } + + public int getAffectedRows() { + return affectedRows; + } + + public Map getData() { + return data; + } + + public Map getOldData() { + return oldData; + } +} \ No newline at end of file diff --git a/flexmodel-server/src/main/java/dev/flexmodel/realtime/RealtimeRabbitmqListener.java b/flexmodel-server/src/main/java/dev/flexmodel/realtime/RealtimeRabbitmqListener.java new file mode 100644 index 00000000..c76098e0 --- /dev/null +++ b/flexmodel-server/src/main/java/dev/flexmodel/realtime/RealtimeRabbitmqListener.java @@ -0,0 +1,106 @@ +package dev.flexmodel.realtime; + +import dev.flexmodel.codegen.entity.Project; +import dev.flexmodel.event.ChangedEvent; +import dev.flexmodel.event.EventListener; +import dev.flexmodel.project.ProjectRepository; +import io.smallrye.reactive.messaging.MutinyEmitter; +import io.smallrye.reactive.messaging.rabbitmq.OutgoingRabbitMQMetadata; +import jakarta.enterprise.context.ApplicationScoped; +import jakarta.enterprise.inject.Instance; +import jakarta.inject.Inject; +import lombok.extern.slf4j.Slf4j; +import org.eclipse.microprofile.reactive.messaging.Channel; +import org.eclipse.microprofile.reactive.messaging.Message; +import org.eclipse.microprofile.reactive.messaging.Metadata; + +import java.util.Map; + +/** + * 数据变更事件 RabbitMQ 监听器,桥接引擎层的 EventPublisher 到 RabbitMQ topic 交换机。 + *

+ * 参照 {@link RealtimeEventListener} 的模式:监听后置事件(INSERTED / UPDATED / DELETED), + * 仅在操作成功时转发。与 WebSocket 实时广播互不影响,各走各的通道。 + *

+ * routing key 形如 {@code data...},投递到与 flow 事件共用的 + * {@code flexmodel.events} topic 交换机,消费端可按项目、模型或操作类型订阅。 + *

+ * 投递为尽力而为:失败仅告警,不重试、不阻塞、不回滚。 + * + * @author cjbi + */ +@Slf4j +@ApplicationScoped +public class RealtimeRabbitmqListener implements EventListener { + + @Inject + @Channel("data-events-out") + Instance> dataEventEmitterInstance; + + @Inject + ProjectRepository projectRepository; + + @Override + public void onChanged(ChangedEvent event) { + // 操作失败时不转发,避免无谓的对象构造与 IO + if (!event.isSuccess()) { + return; + } + if (dataEventEmitterInstance.isUnsatisfied()) { + log.debug("data-events-out channel not resolvable, skip forwarding.||model={}", event.getModelName()); + return; + } + try { + Project project = projectRepository.findProjectByDatabaseName(event.getSchemaName()); + String operation = mapEventType(event.getEventType()); + String routingKey = "data." + project.getId() + "." + event.getModelName() + "." + operation.toLowerCase(); + Map newData = event.getNewData() != null ? event.getNewData() : Map.of(); + Map oldData = event.getOldData() != null ? event.getOldData() : Map.of(); + + DataChangeEvent payload = new DataChangeEvent( + routingKey, project.getId(), operation, event.getModelName(), event.getSchemaName(), + event.getId(), event.getTimestamp(), event.getAffectedRows(), newData, oldData); + + MutinyEmitter emitter = dataEventEmitterInstance.get(); + OutgoingRabbitMQMetadata metadata = new OutgoingRabbitMQMetadata.Builder() + .withRoutingKey(routingKey) + // 持久化投递(delivery_mode=2):使消息在 broker 重启后仍可恢复,配合 durable 交换机与 durable 队列生效 + .withDeliveryMode(2) + .build(); + // 异步发送,不阻塞引擎写入线程;失败仅告警 + emitter.sendMessage(Message.of(payload, Metadata.of(metadata))) + .onFailure() + .invoke(e -> log.warn("forward data change event to rabbitmq failed.||model={}||operation={}", + event.getModelName(), operation, e)) + .subscribe() + .asCompletionStage(); + } catch (Exception e) { + log.warn("forward data change event to rabbitmq failed (sync).||model={}", event.getModelName(), e); + } + } + + @Override + public boolean supports(String eventType) { + return "INSERTED".equals(eventType) + || "UPDATED".equals(eventType) + || "DELETED".equals(eventType); + } + + @Override + public int getOrder() { + // 低于 RealtimeEventListener(1000),不影响业务实时推送 + return 1100; + } + + /** + * 引擎事件类型映射到协议操作类型:INSERTED -> INSERT, UPDATED -> UPDATE, DELETED -> DELETE + */ + private String mapEventType(String engineEventType) { + return switch (engineEventType) { + case "INSERTED" -> "INSERT"; + case "UPDATED" -> "UPDATE"; + case "DELETED" -> "DELETE"; + default -> engineEventType; + }; + } +} diff --git a/flexmodel-server/src/main/java/dev/flexmodel/scheduling/TriggerService.java b/flexmodel-server/src/main/java/dev/flexmodel/scheduling/TriggerService.java index 975eb5f4..25a174a1 100644 --- a/flexmodel-server/src/main/java/dev/flexmodel/scheduling/TriggerService.java +++ b/flexmodel-server/src/main/java/dev/flexmodel/scheduling/TriggerService.java @@ -64,7 +64,7 @@ public class TriggerService { /** * 应用启动时扫描全部启用项目的 f_trigger 表, - * 对其中调度触发器(SCHEDULED)按 scheduler.getTrigger(...) 查询 Quartz 中是否已存在对应调度任务, + * 对其中调度触发器(SCHEDULE)按 scheduler.getTrigger(...) 查询 Quartz 中是否已存在对应调度任务, * 不存在则(重建)调度。用于重启后恢复丢失的调度任务。 *

* 为避免早于项目 Schema 注册执行,这里显式为每个可能尚未注册 Schema 的项目调用 @@ -95,7 +95,7 @@ void restoreScheduledTriggersOnStartup(@Observes StartupEvent event) { /** * 对单个项目同步 f_trigger 中的调度触发器到 Quartz。 - * 遍历该项目 state=true 且 type=SCHEDULED 的触发器, + * 遍历该项目 state=true 且 type=SCHEDULE 的触发器, * 若 scheduler.getTrigger(...) 不存在对应调度任务则按配置重建。 * * @return 本次恢复(新建)的调度任务数量 @@ -104,7 +104,7 @@ private int syncScheduledTriggers(Project project) { String projectId = project.getId(); // 仅处理调度类型触发器;事件触发器(EVENT)不参与 Quartz 调度 List triggers = triggerRepository.find(projectId, - trigger.state.eq(true).and(trigger.type.eq(TriggerType.SCHEDULED)), 1, Integer.MAX_VALUE); + trigger.state.eq(true).and(trigger.type.eq(TriggerType.SCHEDULE)), 1, Integer.MAX_VALUE); int restored = 0; for (Trigger trigger : triggers) { diff --git a/flexmodel-server/src/main/java/dev/flexmodel/storage/config/S3Backend.java b/flexmodel-server/src/main/java/dev/flexmodel/storage/config/S3Backend.java index 9a04be7c..5e2976d5 100644 --- a/flexmodel-server/src/main/java/dev/flexmodel/storage/config/S3Backend.java +++ b/flexmodel-server/src/main/java/dev/flexmodel/storage/config/S3Backend.java @@ -6,6 +6,8 @@ import software.amazon.awssdk.services.s3.S3Client; import software.amazon.awssdk.services.s3.model.*; +import lombok.extern.slf4j.Slf4j; + import java.util.List; import java.util.stream.Collectors; @@ -21,6 +23,7 @@ * * @author cjbi */ +@Slf4j public class S3Backend implements StorageBackend { public static final String TYPE = "s3"; @@ -94,8 +97,34 @@ public boolean containerExists(String prefix) { public void validate() { try { s3Client.headBucket(b -> b.bucket(bucket)); + } catch (NoSuchBucketException e) { + autoCreateBucket(); + } catch (S3Exception e) { + // Some S3-compatible stores (MinIO/RustFS) return 404 for a missing bucket + // but may not map it to NoSuchBucketException, so also match by status code. + if (e.statusCode() == 404) { + autoCreateBucket(); + } else { + throw new InternalServerException("Failed to validate S3 bucket '" + bucket + "': " + e.getMessage(), e); + } + } + } + + /** + * Auto-create the bucket when it does not exist; works for any S3-compatible store. + * In read-only mode the storage is never mutated, so we fail-fast and ask the + * operator to create the bucket manually. + */ + private void autoCreateBucket() { + if (readOnly) { + throw new InternalServerException("S3 bucket '" + bucket + "' does not exist; read-only mode is enabled, refusing to auto-create. Please create the bucket manually."); + } + try { + log.info("S3 bucket '{}' does not exist, auto-creating", bucket); + s3Client.createBucket(b -> b.bucket(bucket)); + log.info("S3 bucket '{}' created successfully", bucket); } catch (S3Exception e) { - throw new InternalServerException("Failed to validate S3 bucket '" + bucket + "': " + e.getMessage(), e); + throw new InternalServerException("Failed to auto-create S3 bucket '" + bucket + "': " + e.getMessage(), e); } } diff --git a/flexmodel-server/src/main/resources/application-dev.properties b/flexmodel-server/src/main/resources/application-dev.properties index 7fbae701..ee35fb3a 100644 --- a/flexmodel-server/src/main/resources/application-dev.properties +++ b/flexmodel-server/src/main/resources/application-dev.properties @@ -1,3 +1,7 @@ quarkus.log.category."io.quarkus.arc".level=INFO quarkus.log.category."io.agroal".level=INFO quarkus.log.category."dev.flexmodel".level=INFO +flexmodel.events.rabbitmq.enabled=true +quarkus.rabbitmq.username=flexmodel +quarkus.rabbitmq.password=flexmodel +quarkus.rabbitmq.devservices.enabled=true diff --git a/flexmodel-server/src/main/resources/application.properties b/flexmodel-server/src/main/resources/application.properties index 36249441..26da7025 100644 --- a/flexmodel-server/src/main/resources/application.properties +++ b/flexmodel-server/src/main/resources/application.properties @@ -69,3 +69,18 @@ flexmodel.jwt.refresh-token-lifetime=30d # flexmodel.storage.s3-bucket=flexmodel # Pages Configuration flexmodel.pages.root-path=./pages +# Flow 生命周期 / 数据变更事件 RabbitMQ 桥接 +# 默认不连 broker;需要 RabbitMQ 转发时置 flexmodel.events.rabbitmq.enabled=true 并提供连接配置。 +# 连接器级默认值(对所有 smallrye-rabbitmq 通道生效) +mp.messaging.connector.smallrye-rabbitmq.exchange.name=flexmodel.events +mp.messaging.connector.smallrye-rabbitmq.exchange.type=topic +mp.messaging.connector.smallrye-rabbitmq.exchange.durable=true +mp.messaging.connector.smallrye-rabbitmq.host=${quarkus.rabbitmq.host:localhost} +mp.messaging.connector.smallrye-rabbitmq.port=${quarkus.rabbitmq.port:5672} +mp.messaging.connector.smallrye-rabbitmq.username=${quarkus.rabbitmq.username:guest} +mp.messaging.connector.smallrye-rabbitmq.password=${quarkus.rabbitmq.password:guest} +mp.messaging.connector.smallrye-rabbitmq.virtual-host=${quarkus.rabbitmq.virtual-host:/} +mp.messaging.outgoing.events-out.connector=smallrye-rabbitmq +mp.messaging.outgoing.events-out.enabled=${flexmodel.events.rabbitmq.enabled:false} +mp.messaging.outgoing.data-events-out.connector=smallrye-rabbitmq +mp.messaging.outgoing.data-events-out.enabled=${flexmodel.events.rabbitmq.enabled:false} diff --git a/flexmodel-server/src/main/resources/dev_test.fml b/flexmodel-server/src/main/resources/dev_test.fml index 472146ea..c17c7cad 100644 --- a/flexmodel-server/src/main/resources/dev_test.fml +++ b/flexmodel-server/src/main/resources/dev_test.fml @@ -103,11 +103,11 @@ seed f_em_flow_deployment { seed f_trigger { {id: "4e8be2da-8ae7-4571-b347-61af5fdd3e9b", name: "请求日志归档任务", description: "写入日志归档表进行数据归档", type: "EVENT", config: {type: "event", modelName: "api_request_log", mutationTypes: ["create"], triggerTiming: "after", datasourceName: "system"}, job_type: "FLOW", job_group: "system_api_request_log", job_id: "22203aa0-5f25-4b8a-974d-e4288c1bbd18", state: false, created_at: "2025-10-19 17:59:45", updated_at: "2025-10-19 20:38:24", created_by: null, updated_by: null}, - {id: "bf492f37-1f01-4eb8-b76d-d319299b4d8e", name: "定时触发-间隔触发", description: "定时触发-间隔触发-备注", type: "SCHEDULED", config: {type: "interval", interval: 1, intervalUnit: "minute", repeatCount: 100}, job_id: "5c41f37a-87a9-47af-bdba-44d0c27eda89", job_type: "FLOW", job_group: "testFlowName_1757909108656", state: false, created_at: "2025-09-15 12:05:09.000", updated_at: "2025-09-15 12:05:09.000"}, - {id: "d8c60d2a-19d8-4c3c-b370-96318733858f", name: "定时触发-Cron表达式", description: "定时触发-Cron表达式-备注", type: "SCHEDULED", config: {type: "cron", cronExpression: "0 0 * * * ? *"}, job_id: "5c41f37a-87a9-47af-bdba-44d0c27eda89", job_type: "FLOW", job_group: "testFlowName_1757909108656", state: false, created_at: "2025-09-15 12:05:09.000", updated_at: "2025-09-15 12:05:09.000"}, + {id: "bf492f37-1f01-4eb8-b76d-d319299b4d8e", name: "定时触发-间隔触发", description: "定时触发-间隔触发-备注", type: "SCHEDULE", config: {type: "interval", interval: 1, intervalUnit: "minute", repeatCount: 100}, job_id: "5c41f37a-87a9-47af-bdba-44d0c27eda89", job_type: "FLOW", job_group: "testFlowName_1757909108656", state: false, created_at: "2025-09-15 12:05:09.000", updated_at: "2025-09-15 12:05:09.000"}, + {id: "d8c60d2a-19d8-4c3c-b370-96318733858f", name: "定时触发-Cron表达式", description: "定时触发-Cron表达式-备注", type: "SCHEDULE", config: {type: "cron", cronExpression: "0 0 * * * ? *"}, job_id: "5c41f37a-87a9-47af-bdba-44d0c27eda89", job_type: "FLOW", job_group: "testFlowName_1757909108656", state: false, created_at: "2025-09-15 12:05:09.000", updated_at: "2025-09-15 12:05:09.000"}, {id: "9434666d-6b3b-417a-80e3-3352f40ded71", name: "事件触发-之前", description: "事件触发-之前-描述", type: "EVENT", config: {type: "event", datasourceName: "dev_test", modelName: "Student", mutationTypes: ["delete", "update", "create"], triggerTiming: "after"}, job_id: "5c41f37a-87a9-47af-bdba-44d0c27eda89", job_type: "FLOW", job_group: "dev_test_Student", state: true, created_at: "2025-09-15 12:05:09.000", updated_at: "2025-09-15 12:05:09.000"}, {id: "f351b8d9-a450-4f2c-8fff-cf862d690352", name: "事件触发-之后", description: "事件触发-之后-描述", type: "EVENT", config: {type: "event", datasourceName: "dev_test", modelName: "Classes", mutationTypes: ["create", "update", "delete"], triggerTiming: "after"}, job_id: "5c41f37a-87a9-47af-bdba-44d0c27eda89", job_type: "FLOW", job_group: "dev_test_Classes", state: true, created_at: "2025-09-15 12:05:09.000", updated_at: "2025-09-15 12:05:09.000"}, - {id: "687d14a2-9d5e-4ecf-b85c-dfc7c3b31493", name: "函数运行时健康检测", description: null, type: "SCHEDULED", config: {type: "interval", interval: 10, intervalUnit: "second", repeatCount: 999999}, job_type: "FUNCTION", job_group: "dev_test_fn_health-check", job_id: "health-check", state: true, created_at: "2026-07-15T10:50:31.284785500", updated_at: "2026-07-17T12:02:48.970596400", created_by: null, updated_by: null}, + {id: "687d14a2-9d5e-4ecf-b85c-dfc7c3b31493", name: "函数运行时健康检测", description: null, type: "SCHEDULE", config: {type: "interval", interval: 10, intervalUnit: "second", repeatCount: 999999}, job_type: "FUNCTION", job_group: "dev_test_fn_health-check", job_id: "health-check", state: true, created_at: "2026-07-15T10:50:31.284785500", updated_at: "2026-07-17T12:02:48.970596400", created_by: null, updated_by: null}, } seed f_function { diff --git a/flexmodel-server/src/main/resources/project.fml b/flexmodel-server/src/main/resources/project.fml index 7d531077..58addb24 100644 --- a/flexmodel-server/src/main/resources/project.fml +++ b/flexmodel-server/src/main/resources/project.fml @@ -285,7 +285,7 @@ model f_function_log { enum TriggerType { EVENT, - SCHEDULED, + SCHEDULE, @system } diff --git a/flexmodel-server/src/test/java/dev/flexmodel/flow/event/FlowEventPublisherTest.java b/flexmodel-server/src/test/java/dev/flexmodel/flow/event/FlowEventPublisherTest.java new file mode 100644 index 00000000..40975a57 --- /dev/null +++ b/flexmodel-server/src/test/java/dev/flexmodel/flow/event/FlowEventPublisherTest.java @@ -0,0 +1,88 @@ +package dev.flexmodel.flow.event; + +import dev.flexmodel.SQLiteTestResource; +import io.quarkus.test.common.QuarkusTestResource; +import io.quarkus.test.junit.QuarkusTest; +import jakarta.inject.Inject; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * 验证 {@link FlowEventPublisher} 将强类型事件经 Vert.x EventBus 广播, + * 且事件字段完整保留。本地事件为机制本身,不依赖 RabbitMQ。 + * + * @author cjbi + */ +@QuarkusTest +@QuarkusTestResource(SQLiteTestResource.class) +public class FlowEventPublisherTest { + + @Inject + FlowEventPublisher flowEventPublisher; + + @Inject + FlowEventTestConsumer capturingConsumer; + + @Test + void publishFlowInstanceStartedEventFieldsPreserved() { + capturingConsumer.captured().set(null); + + String projectId = "proj-test"; + String caller = "caller-test"; + String flowDeployId = "deploy-test"; + String flowInstanceId = "instance-test"; + Map variables = new HashMap<>(); + variables.put("amount", 100); + variables.put("approved", true); + + flowEventPublisher.publish(new FlowInstanceStartedEvent(projectId, caller, flowDeployId, + flowInstanceId, variables)); + + FlowInstanceStartedEvent event = pollFor(capturingConsumer.captured(), 5000L); + + assertNotNull(event, "FlowInstanceStartedEvent should be received via EventBus"); + assertEquals(projectId, event.getProjectId()); + assertEquals(caller, event.getCaller()); + assertEquals(flowDeployId, event.getFlowDeployId()); + assertEquals(flowInstanceId, event.getFlowInstanceId()); + assertNotNull(event.getVariables()); + assertEquals(100, event.getVariables().get("amount")); + assertEquals(Boolean.TRUE, event.getVariables().get("approved")); + assertEquals(FlowEventTypes.FLOW_INSTANCE_STARTED, event.routingKey()); + assertTrue(event.getTimestamp() > 0, "timestamp should be set at construction"); + } + + @Test + void publishNullEventIsNoOp() { + // 不应抛出,也不应影响其他消费者 + flowEventPublisher.publish(null); + } + + /** + * 有界轮询等待消费者捕获事件,避免引入 Awaitility 依赖。 + */ + private static T pollFor(AtomicReference ref, long timeoutMs) { + long deadline = System.currentTimeMillis() + timeoutMs; + T value = ref.get(); + while (value == null && System.currentTimeMillis() < deadline) { + sleep(50L); + value = ref.get(); + } + return value; + } + + private static void sleep(long millis) { + try { + Thread.sleep(millis); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } +} diff --git a/flexmodel-server/src/test/java/dev/flexmodel/flow/event/FlowEventRabbitmqBridgeE2ETest.java b/flexmodel-server/src/test/java/dev/flexmodel/flow/event/FlowEventRabbitmqBridgeE2ETest.java new file mode 100644 index 00000000..2607610b --- /dev/null +++ b/flexmodel-server/src/test/java/dev/flexmodel/flow/event/FlowEventRabbitmqBridgeE2ETest.java @@ -0,0 +1,174 @@ +package dev.flexmodel.flow.event; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.rabbitmq.client.BuiltinExchangeType; +import com.rabbitmq.client.Channel; +import com.rabbitmq.client.Connection; +import com.rabbitmq.client.ConnectionFactory; +import com.rabbitmq.client.Envelope; +import com.rabbitmq.client.GetResponse; +import dev.flexmodel.SQLiteTestResource; +import io.quarkus.test.common.QuarkusTestResource; +import io.quarkus.test.junit.QuarkusTest; +import jakarta.inject.Inject; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.condition.EnabledIfEnvironmentVariable; + +import java.util.HashMap; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * 端到端验证 {@link FlowEventRabbitmqBridge} 将 Flow 生命周期事件以正确 routing key 与 + * JSON 载荷推送到 RabbitMQ topic 交换机。 + *

+ * 使用 Testcontainers 启动真实 RabbitMQ broker(见 {@link RabbitMqTestResource}), + * 并以 amqp-client 临时队列订阅交换机、拉取转发的消息。 + *

+ * 默认跳过(避免无 Docker 环境下默认 {@code mvn test} 失败)。运行需: + *

    + *
  • 本机 Docker 可用且能拉取 {@code rabbitmq:3-management} 镜像
  • + *
  • 设置环境变量 {@code FLEXMODEL_E2E_RABBITMQ=true}(opt-in 启用本端到端测试)
  • + *
+ *

+ * {@code @QuarkusTestResource} 设 {@code restrictToAnnotatedClass=true},确保 broker 资源 + * 仅对本测试类生效,不泄漏到其他共享应用上下文的测试。 + * + * @author cjbi + */ +@EnabledIfEnvironmentVariable(named = "FLEXMODEL_E2E_RABBITMQ", matches = "true") +@QuarkusTest +@QuarkusTestResource(SQLiteTestResource.class) +@QuarkusTestResource(value = RabbitMqTestResource.class, restrictToAnnotatedClass = true) +public class FlowEventRabbitmqBridgeE2ETest { + + private static final String EXCHANGE = "flexmodel.events"; + private static final ObjectMapper MAPPER = new ObjectMapper(); + + @Inject + FlowEventPublisher flowEventPublisher; + + /** + * 验证带 variables 快照的 FlowInstanceStartedEvent 经桥接转发后,variables 仍在 JSON 载荷中完整。 + */ + @Test + void flowInstanceStartedEventWithVariablesForwarded() throws Exception { + String projectId = "proj-vars"; + String caller = "caller-vars"; + String flowDeployId = "deploy-vars"; + String flowInstanceId = "instance-vars"; + Map variables = new HashMap<>(); + variables.put("amount", 100); + variables.put("approved", true); + + try (Connection connection = newConnection(); + Channel channel = connection.createChannel()) { + channel.exchangeDeclare(EXCHANGE, BuiltinExchangeType.TOPIC, true); + String queue = channel.queueDeclare().getQueue(); + channel.queueBind(queue, EXCHANGE, "flow.proj-vars.instance.started"); + + flowEventPublisher.publish( + new FlowInstanceStartedEvent(projectId, caller, flowDeployId, flowInstanceId, variables)); + + GetResponse response = pollForMessage(channel, queue, 15000L); + assertNotNull(response, "FlowInstanceStartedEvent should be forwarded to broker within timeout"); + + assertEquals("flow.proj-vars.instance.started", response.getEnvelope().getRoutingKey()); + + JsonNode body = MAPPER.readTree(response.getBody()); + assertEquals(projectId, body.get("projectId").asText()); + assertEquals(flowDeployId, body.get("flowDeployId").asText()); + assertEquals(flowInstanceId, body.get("flowInstanceId").asText()); + JsonNode vars = body.get("variables"); + assertNotNull(vars, "variables snapshot should be present in payload"); + assertEquals(100, vars.get("amount").asInt()); + assertTrue(vars.get("approved").asBoolean(), "approved variable should be true"); + } + } + + /** + * 验证 UserTaskSuspendedEvent 经桥接转发后,routing key 正确,且 nodeAttributes(节点定义扩展属性快照) + * 与 variables 一样在 JSON 载荷中完整保留,外部订阅者无需回查定义仓库即可读取节点配置。 + */ + @Test + void userTaskSuspendedEventWithNodeAttributesForwarded() throws Exception { + String projectId = "proj-task"; + String caller = "caller-task"; + String flowDeployId = "deploy-task"; + String flowInstanceId = "instance-task"; + String nodeInstanceId = "nodeinst-task"; + String nodeKey = "approveNode"; + Map variables = new HashMap<>(); + variables.put("amount", 500); + Map nodeAttributes = new HashMap<>(); + nodeAttributes.put("name", "审批"); + nodeAttributes.put("assignee", "manager"); + nodeAttributes.put("multiInstance", false); + + try (Connection connection = newConnection(); + Channel channel = connection.createChannel()) { + channel.exchangeDeclare(EXCHANGE, BuiltinExchangeType.TOPIC, true); + String queue = channel.queueDeclare().getQueue(); + channel.queueBind(queue, EXCHANGE, "flow.proj-task.usertask.suspended"); + + flowEventPublisher.publish(new UserTaskSuspendedEvent(projectId, caller, flowDeployId, + flowInstanceId, nodeInstanceId, nodeKey, variables, nodeAttributes)); + + GetResponse response = pollForMessage(channel, queue, 15000L); + assertNotNull(response, "UserTaskSuspendedEvent should be forwarded to broker within timeout"); + + assertEquals("flow.proj-task.usertask.suspended", response.getEnvelope().getRoutingKey()); + + JsonNode body = MAPPER.readTree(response.getBody()); + assertEquals(projectId, body.get("projectId").asText()); + assertEquals(flowDeployId, body.get("flowDeployId").asText()); + assertEquals(flowInstanceId, body.get("flowInstanceId").asText()); + assertEquals(nodeInstanceId, body.get("nodeInstanceId").asText()); + assertEquals(nodeKey, body.get("nodeKey").asText()); + JsonNode vars = body.get("variables"); + assertNotNull(vars, "variables snapshot should be present"); + assertEquals(500, vars.get("amount").asInt()); + JsonNode attrs = body.get("nodeAttributes"); + assertNotNull(attrs, "nodeAttributes snapshot should be present in payload"); + assertEquals("审批", attrs.get("name").asText()); + assertEquals("manager", attrs.get("assignee").asText()); + assertFalse(attrs.get("multiInstance").asBoolean(), "multiInstance should be false"); + } + } + + private static Connection newConnection() throws Exception { + ConnectionFactory factory = new ConnectionFactory(); + factory.setHost(RabbitMqTestResource.host); + factory.setPort(RabbitMqTestResource.port); + factory.setUsername(RabbitMqTestResource.username); + factory.setPassword(RabbitMqTestResource.password); + return factory.newConnection(); + } + + /** + * 有界轮询拉取消息:桥接转发为异步,需留足时间。 + */ + private static GetResponse pollForMessage(Channel channel, String queue, long timeoutMs) + throws Exception { + long deadline = System.currentTimeMillis() + timeoutMs; + GetResponse response = channel.basicGet(queue, true); + while (response == null && System.currentTimeMillis() < deadline) { + sleep(100L); + response = channel.basicGet(queue, true); + } + return response; + } + + private static void sleep(long millis) { + try { + Thread.sleep(millis); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } +} diff --git a/flexmodel-server/src/test/java/dev/flexmodel/flow/event/FlowEventRabbitmqBridgeTest.java b/flexmodel-server/src/test/java/dev/flexmodel/flow/event/FlowEventRabbitmqBridgeTest.java new file mode 100644 index 00000000..07ceb640 --- /dev/null +++ b/flexmodel-server/src/test/java/dev/flexmodel/flow/event/FlowEventRabbitmqBridgeTest.java @@ -0,0 +1,83 @@ +package dev.flexmodel.flow.event; + +import dev.flexmodel.SQLiteTestResource; +import io.quarkus.test.common.QuarkusTestResource; +import io.quarkus.test.junit.QuarkusTest; +import jakarta.inject.Inject; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * 验证 FlowEventRabbitmqBridge 的可选桥接语义:默认通道禁用且不连 broker 时, + * 桥接 bean 仍可创建,本地事件照常经 FlowEventPublisher 广播。 + * 不依赖真实 broker,亦不依赖 SmallRye InMemoryConnector(当前 Quarkus 版本未提供配套扩展), + * 仅以默认配置启动应用即可断言桥接的零侵入设计目标。 + * + * @author cjbi + */ +@QuarkusTest +@QuarkusTestResource(SQLiteTestResource.class) +public class FlowEventRabbitmqBridgeTest { + + @Inject + FlowEventRabbitmqBridge bridge; + + @Inject + FlowEventPublisher flowEventPublisher; + + @Inject + FlowEventTestConsumer capturingConsumer; + + @Test + void bridgeBeanPresentAndNonThrowingByDefault() { + // 默认通道禁用:桥接 bean 仍存在(零侵入),发布桥接也消费的事件不应抛出, + // 且不连 broker(应用无 broker 启动即证) + assertNotNull(bridge, "bridge bean should be present (zero-intrusion)"); + Map variables = new HashMap<>(); + variables.put("k", "v"); + flowEventPublisher.publish(new FlowInstanceStartedEvent("proj-bridge", "caller-bridge", + "deploy-bridge", "instance-bridge", variables)); + } + + @Test + void disabledBridgeIsTransparentToLocalEvents() { + capturingConsumer.captured().set(null); + + Map variables = new HashMap<>(); + variables.put("k", "v"); + flowEventPublisher.publish(new FlowInstanceStartedEvent("proj-bridge", "caller-bridge", + "deploy-bridge", "instance-bridge", variables)); + + FlowInstanceStartedEvent event = pollFor(capturingConsumer.captured(), 5000L); + + assertNotNull(event, "local EventBus consumer should still receive event while channel disabled"); + assertEquals("proj-bridge", event.getProjectId()); + assertEquals("instance-bridge", event.getFlowInstanceId()); + assertTrue(event.getTimestamp() > 0); + } + + private static T pollFor(AtomicReference ref, long timeoutMs) { + long deadline = System.currentTimeMillis() + timeoutMs; + T value = ref.get(); + while (value == null && System.currentTimeMillis() < deadline) { + sleep(50L); + value = ref.get(); + } + return value; + } + + private static void sleep(long millis) { + try { + Thread.sleep(millis); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } +} diff --git a/flexmodel-server/src/test/java/dev/flexmodel/flow/event/FlowEventTestConsumer.java b/flexmodel-server/src/test/java/dev/flexmodel/flow/event/FlowEventTestConsumer.java new file mode 100644 index 00000000..ab7d628f --- /dev/null +++ b/flexmodel-server/src/test/java/dev/flexmodel/flow/event/FlowEventTestConsumer.java @@ -0,0 +1,28 @@ +package dev.flexmodel.flow.event; + +import io.quarkus.vertx.ConsumeEvent; +import jakarta.enterprise.context.ApplicationScoped; + +import java.util.concurrent.atomic.AtomicReference; + +/** + * 测试用 EventBus 消费者:捕获 {@link FlowInstanceStartedEvent},验证字段经 EventBus 透传完整。 + *

+ * 顶层 {@code @ApplicationScoped} Bean,确保 Quarkus 发现并注册其 {@code @ConsumeEvent} 消费者。 + * + * @author cjbi + */ +@ApplicationScoped +public class FlowEventTestConsumer { + + final AtomicReference captured = new AtomicReference<>(); + + @ConsumeEvent(value = FlowEventTypes.FLOW_INSTANCE_STARTED, blocking = false) + void onFlowInstanceStarted(FlowInstanceStartedEvent event) { + captured.set(event); + } + + public AtomicReference captured() { + return captured; + } +} diff --git a/flexmodel-server/src/test/java/dev/flexmodel/flow/event/RabbitMqTestResource.java b/flexmodel-server/src/test/java/dev/flexmodel/flow/event/RabbitMqTestResource.java new file mode 100644 index 00000000..e926ddd4 --- /dev/null +++ b/flexmodel-server/src/test/java/dev/flexmodel/flow/event/RabbitMqTestResource.java @@ -0,0 +1,59 @@ +package dev.flexmodel.flow.event; + +import io.quarkus.test.common.QuarkusTestResourceLifecycleManager; +import org.testcontainers.containers.RabbitMQContainer; + +import java.util.LinkedHashMap; +import java.util.Map; + +/** + * Testcontainers 启动 RabbitMQ broker,并向 Quarkus 注入连接配置与启用 events-out 通道。 + *

+ * 用于端到端验证 FlowEventRabbitmqBridge 将事件以正确 routing key 与 JSON 载荷推送到 topic 交换机。 + *

+ * 连接配置注入到 SmallRye RabbitMQ connector 的 channel 级属性 + * ({@code mp.messaging.outgoing.events-out.host/port/username/password})。 + * 注意:{@code quarkus-messaging-rabbitmq} 扩展的 {@code quarkus.rabbitmq.*} 只注册 DevServices 与 + * credentials provider 字段,不注册 host/port/username/password 连接字段——这些由 SmallRye + * connector 自身读取(channel 级或全局 {@code rabbitmq-*} 别名)。显式提供 channel 级连接配置后, + * DevServices 自动旁路。 + * + * @author cjbi + */ +public class RabbitMqTestResource implements QuarkusTestResourceLifecycleManager { + + // 静态引用,供测试用例获取连接坐标(测试侧 amqp-client 用其订阅交换机) + static volatile String host; + static volatile int port; + static volatile String username = "guest"; + static volatile String password = "guest"; + + private RabbitMQContainer rabbitmq; + + @Override + public Map start() { + rabbitmq = new RabbitMQContainer("rabbitmq:3-management"); + rabbitmq.start(); + host = rabbitmq.getHost(); + port = rabbitmq.getAmqpPort(); + + Map config = new LinkedHashMap<>(); + // SmallRye RabbitMQ connector channel 级连接配置(quarkus.rabbitmq.* 不注册连接字段) + config.put("mp.messaging.outgoing.events-out.host", host); + config.put("mp.messaging.outgoing.events-out.port", String.valueOf(port)); + config.put("mp.messaging.outgoing.events-out.username", username); + config.put("mp.messaging.outgoing.events-out.password", password); + // 显式提供连接配置后 DevServices 自动旁路;此处再确认禁用,避免无 Docker 时误触 + config.put("quarkus.rabbitmq.devservices.enabled", "false"); + // 启用 events-out 通道,使桥接真正转发到 broker + config.put("mp.messaging.outgoing.events-out.enabled", "true"); + return config; + } + + @Override + public void stop() { + if (rabbitmq != null) { + rabbitmq.close(); + } + } +} diff --git a/flexmodel-server/src/test/java/dev/flexmodel/pages/PageAliasManagerTest.java b/flexmodel-server/src/test/java/dev/flexmodel/pages/PageAliasManagerTest.java index 41283e55..fa499fe3 100644 --- a/flexmodel-server/src/test/java/dev/flexmodel/pages/PageAliasManagerTest.java +++ b/flexmodel-server/src/test/java/dev/flexmodel/pages/PageAliasManagerTest.java @@ -206,6 +206,11 @@ public PagesConfig pages() { return () -> pagesRootPath; } + @Override + public EventsConfig events() { + return null; + } + @Override public String projectUrlTemplate() { return ""; } diff --git a/flexmodel-server/src/test/java/dev/flexmodel/pages/PageDeployerTest.java b/flexmodel-server/src/test/java/dev/flexmodel/pages/PageDeployerTest.java index 3bb52c40..209c5409 100644 --- a/flexmodel-server/src/test/java/dev/flexmodel/pages/PageDeployerTest.java +++ b/flexmodel-server/src/test/java/dev/flexmodel/pages/PageDeployerTest.java @@ -235,6 +235,11 @@ public PagesConfig pages() { return () -> pagesRootPath; } + @Override + public EventsConfig events() { + return null; + } + @Override public String projectUrlTemplate() { return ""; } diff --git a/flexmodel-server/src/test/java/dev/flexmodel/rest/PermissionFilterTest.java b/flexmodel-server/src/test/java/dev/flexmodel/rest/PermissionFilterTest.java index b910f0c0..d72ff4fb 100644 --- a/flexmodel-server/src/test/java/dev/flexmodel/rest/PermissionFilterTest.java +++ b/flexmodel-server/src/test/java/dev/flexmodel/rest/PermissionFilterTest.java @@ -393,7 +393,7 @@ void schedulingExecuteShouldPassCreateTrigger() { .header("Authorization", testTokenHelper.getAuthorizationHeader()) .header("X-Test-Permissions", "scheduling:execute") .contentType(ContentType.JSON) - .body("{\"name\":\"TestTrigger\",\"type\":\"SCHEDULED\"}") + .body("{\"name\":\"TestTrigger\",\"type\":\"SCHEDULE\"}") .when() .post(Resources.ROOT_PATH + "/projects/dev_test/triggers") .then() diff --git a/flexmodel-server/src/test/java/dev/flexmodel/rest/TriggerResourceTest.java b/flexmodel-server/src/test/java/dev/flexmodel/rest/TriggerResourceTest.java index 614ae3cb..17f74e3a 100644 --- a/flexmodel-server/src/test/java/dev/flexmodel/rest/TriggerResourceTest.java +++ b/flexmodel-server/src/test/java/dev/flexmodel/rest/TriggerResourceTest.java @@ -45,14 +45,14 @@ public class TriggerResourceTest { /** * 与 {@code TriggerService.getJobGroup} 推导规则保持一致: - * SCHEDULED + FLOW → "{projectId}_flow_{jobId}";SCHEDULED + FUNCTION → "{projectId}_fn_{jobId}"。 + * SCHEDULE + FLOW → "{projectId}_flow_{jobId}";SCHEDULE + FUNCTION → "{projectId}_fn_{jobId}"。 */ private static final String EXPECTED_FLOW_JOB_GROUP = PROJECT_ID + "_flow_" + TEST_JOB_ID; private String createdTriggerId; /** - * 本用例在测试过程中被显式启用(PATCH state=true)的 seed SCHEDULED 触发器 ID 列表, + * 本用例在测试过程中被显式启用(PATCH state=true)的 seed SCHEDULE 触发器 ID 列表, * 用于在 @AfterEach 中确保其 Quartz 调度任务被取消并还原为 state=false,避免污染后续用例。 */ private final List enabledSeedTriggerIds = new ArrayList<>(); @@ -82,7 +82,7 @@ void tearDown() { .body("{ \"state\": false }") .when() .patch(BASE_PATH + "/" + id); - // seed SCHEDULED FLOW 触发器的 jobGroup 在 update 后会被重算为标准格式 + // seed SCHEDULE FLOW 触发器的 jobGroup 在 update 后会被重算为标准格式 unscheduleFromScheduler(id, EXPECTED_FLOW_JOB_GROUP); // 同时兼容 seed 中旧格式的 jobGroup,避免遗留调度任务 unscheduleFromScheduler(id, "testFlowName_1757909108656"); @@ -125,7 +125,7 @@ private void unscheduleFromScheduler(String triggerId, String jobGroup) { } /** - * 断言 SCHEDULED 触发器已正确调度到 Quartz:JobDetail 与 Trigger 同时存在, + * 断言 SCHEDULE 触发器已正确调度到 Quartz:JobDetail 与 Trigger 同时存在, * 触发器状态为 NORMAL,且 JobDataMap 携带 triggerId/jobId/projectId。 *

* 注:SimpleTrigger 在短间隔下可能在断言前已触发完成并被移除,nextFireTime 可能为 null, @@ -154,13 +154,13 @@ private void assertScheduledInQuartz(String triggerId, String jobGroup) throws S } /** - * 断言触发器未在 Quartz 中创建作业/触发器(用于 state=false 的 SCHEDULED 与 EVENT 类型)。 + * 断言触发器未在 Quartz 中创建作业/触发器(用于 state=false 的 SCHEDULE 与 EVENT 类型)。 */ private void assertNotScheduledInQuartz(String triggerId, String jobGroup) throws SchedulerException { assertFalse(quartz.checkExists(jobKey(triggerId, jobGroup)), - "state=false 的 SCHEDULED 触发器不应创建 Quartz JobDetail: " + jobKey(triggerId, jobGroup)); + "state=false 的 SCHEDULE 触发器不应创建 Quartz JobDetail: " + jobKey(triggerId, jobGroup)); assertFalse(quartz.checkExists(triggerKey(triggerId, jobGroup)), - "state=false 的 SCHEDULED 触发器不应创建 Quartz Trigger: " + triggerKey(triggerId, jobGroup)); + "state=false 的 SCHEDULE 触发器不应创建 Quartz Trigger: " + triggerKey(triggerId, jobGroup)); } /** @@ -224,7 +224,7 @@ void testCreateIntervalTrigger() throws SchedulerException { { "name": "测试间隔触发", "description": "测试描述", - "type": "SCHEDULED", + "type": "SCHEDULE", "config": { "type": "interval", "interval": 5, @@ -248,7 +248,7 @@ void testCreateIntervalTrigger() throws SchedulerException { .body("id", notNullValue()) .body("name", equalTo("测试间隔触发")) .body("description", equalTo("测试描述")) - .body("type", equalTo("SCHEDULED")) + .body("type", equalTo("SCHEDULE")) .body("state", equalTo(true)) .body("jobId", equalTo(TEST_JOB_ID)) .body("jobType", equalTo("FLOW")) @@ -273,7 +273,7 @@ void testCreateCronTrigger() throws SchedulerException { { "name": "测试Cron触发", "description": "测试Cron描述", - "type": "SCHEDULED", + "type": "SCHEDULE", "config": { "type": "cron", "cronExpression": "0 0 8 * * ?" @@ -294,7 +294,7 @@ void testCreateCronTrigger() throws SchedulerException { .statusCode(200) .body("id", notNullValue()) .body("name", equalTo("测试Cron触发")) - .body("type", equalTo("SCHEDULED")) + .body("type", equalTo("SCHEDULE")) .body("jobGroup", equalTo(EXPECTED_FLOW_JOB_GROUP)) .body("config.cronExpression", equalTo("0 0 8 * * ?")) .extract() @@ -362,7 +362,7 @@ void testCreateDisabledTriggerNotScheduled() throws SchedulerException { String triggerJson = """ { "name": "禁用的间隔触发", - "type": "SCHEDULED", + "type": "SCHEDULE", "config": { "type": "interval", "interval": 1, @@ -398,7 +398,7 @@ void testCreateDisabledTriggerNotScheduled() throws SchedulerException { .statusCode(200) .body("state", equalTo(false)); - // state=false 的 SCHEDULED 触发器不应在 Quartz 中创建调度任务 + // state=false 的 SCHEDULE 触发器不应在 Quartz 中创建调度任务 assertNotScheduledInQuartz(triggerId, EXPECTED_FLOW_JOB_GROUP); } @@ -417,7 +417,7 @@ void testUpdateTrigger() throws SchedulerException { "id": "%s", "name": "更新后的触发器名称", "description": "更新后的描述", - "type": "SCHEDULED", + "type": "SCHEDULE", "config": { "type": "interval", "interval": 10, @@ -455,7 +455,7 @@ void testUpdateTrigger() throws SchedulerException { "id": "%s", "name": "定时触发-间隔触发", "description": "定时触发-间隔触发-备注", - "type": "SCHEDULED", + "type": "SCHEDULE", "config": { "type": "interval", "interval": 1, @@ -508,7 +508,7 @@ void testUpdateTriggerReschedules() throws SchedulerException { { "id": "%s", "name": "定时触发-间隔触发", - "type": "SCHEDULED", + "type": "SCHEDULE", "config": { "type": "interval", "interval": 30, @@ -543,7 +543,7 @@ void testUpdateTriggerReschedules() throws SchedulerException { "id": "%s", "name": "定时触发-间隔触发", "description": "定时触发-间隔触发-备注", - "type": "SCHEDULED", + "type": "SCHEDULE", "config": { "type": "interval", "interval": 1, @@ -571,7 +571,7 @@ void testUpdateTriggerReschedules() throws SchedulerException { /** * 测试部分更新触发器 - 只更新状态为false *

- * CRON_TRIGGER_ID 是 seed 中 state=false 的 SCHEDULED Cron 触发器,本用例 PATCH 保持 false, + * CRON_TRIGGER_ID 是 seed 中 state=false 的 SCHEDULE Cron 触发器,本用例 PATCH 保持 false, * 不应触发调度。补充断言:Quartz 中不应存在其调度任务。 */ @Test @@ -597,12 +597,12 @@ void testPatchTriggerDisable() throws SchedulerException { /** * 测试部分更新触发器 - 启用触发器(state: false → true) *

- * INTERVAL_TRIGGER_ID 是 seed 中 state=false 的 SCHEDULED 间隔触发器, + * INTERVAL_TRIGGER_ID 是 seed 中 state=false 的 SCHEDULE 间隔触发器, * PATCH state=true → TriggerService 会依据配置调度到 Quartz;PATCH 回 false → 取消调度。 */ @Test void testPatchTriggerEnable() throws SchedulerException { - // INTERVAL_TRIGGER_ID 是seed中state=false的SCHEDULED触发器 + // INTERVAL_TRIGGER_ID 是seed中state=false的SCHEDULE触发器 // PATCH state=true → TriggerService会调度到Quartz String patchJson = """ { "state": true } @@ -700,7 +700,7 @@ void testDeleteTrigger() throws SchedulerException { { "name": "待删除的触发器", "description": "用于删除测试", - "type": "SCHEDULED", + "type": "SCHEDULE", "config": { "type": "interval", "interval": 1, @@ -762,7 +762,7 @@ void testValidateDifferentTriggerTypes() { .get(BASE_PATH + "/" + INTERVAL_TRIGGER_ID) .then() .statusCode(200) - .body("type", equalTo("SCHEDULED")) + .body("type", equalTo("SCHEDULE")) .body("config.interval", notNullValue()) .body("config.intervalUnit", notNullValue()) .body("config.repeatCount", notNullValue()); @@ -774,7 +774,7 @@ void testValidateDifferentTriggerTypes() { .get(BASE_PATH + "/" + CRON_TRIGGER_ID) .then() .statusCode(200) - .body("type", equalTo("SCHEDULED")) + .body("type", equalTo("SCHEDULE")) .body("config.cronExpression", equalTo("0 0 * * * ? *")); // 验证事件触发配置 diff --git a/flexmodel-ui b/flexmodel-ui index ca8088c7..56ef2736 160000 --- a/flexmodel-ui +++ b/flexmodel-ui @@ -1 +1 @@ -Subproject commit ca8088c7e2c476ee307a83d61cde178c3528675a +Subproject commit 56ef273607971f21673f418d026fa1bbc5e27eed diff --git a/flexmodel-website b/flexmodel-website index c6a5fadb..f180b751 160000 --- a/flexmodel-website +++ b/flexmodel-website @@ -1 +1 @@ -Subproject commit c6a5fadbf9be52ea2a5cc32f4aff6f00470a8fba +Subproject commit f180b751c69c0b6941c6f4cddcd5170fc65e28e4 diff --git a/progress.md b/progress.md index 6cd2030e..89b87444 100644 --- a/progress.md +++ b/progress.md @@ -298,3 +298,180 @@ Quartz 作业被正确创建/移除/状态变更。 - 所有 `@QuarkusTest` 测试因 `SessionFactory` 在测试环境不可用而启动失败 - LSP 错误 (codegen 相关类型如 FlowDefinition 在 IDE 中未解析) — 非 our edits 引起 + +## Flow 生命周期事件:本地 EventBus 事件 + 可选 RabbitMQ 桥接(2026-08-19) + +**目标:** 在 flow 核心关键生命周期点发布强类型本地事件,经 Vert.x EventBus 广播;新增可选 RabbitMQ 桥接,默认关闭、未配置不连 +broker,本地事件照常发布。 + +**实现:** + +- 新增 `dev.flexmodel.flow.event` 包:`FlowEvent` 抽象基类(projectId/caller/timestamp + 抽象 `routingKey()`)、 + `FlowEventTypes` 路由 key 常量、11 个具体事件类(定义层 4 / 实例层 4 / 用户任务层 3)。 +- `FlowEventPublisher`(`@ApplicationScoped`):注入 `EventBus`,`publish(FlowEvent)` 经 + `eventBus.publish(routingKey, event)` 广播,全程 try/catch 仅告警、不抛出、不阻塞流程。 +- `FlowEventRabbitmqBridge`(`@ApplicationScoped`):每事件一个 `@ConsumeEvent(常量, blocking=false)` 方法消费具体类型, + `enabled` 默认 false 时 early-return;`forward()` 以 `Instance` 延迟注入(禁用通道时 bean 仍可创建,不影响本地事件),用 + `OutgoingRabbitMQMetadata` 设 routing key 尽力转发。 +- `FlowEventConfig` `@ConfigMapping(prefix="flexmodel.flow")` 声明 `events.rabbitmq.enabled` 配置根,避免 SmallRye 严格校验报 + `does not map to any root`。 +- 依赖:`flexmodel-server/pom.xml` 加 `io.quarkus:quarkus-messaging-rabbitmq`(3.33 起更名自 + `quarkus-smallrye-reactive-messaging-rabbitmq`,版本由 quarkus-bom 管理)。 +- 配置 `application.properties`:`flexmodel.flow.events.rabbitmq.enabled=false`、`mp.messaging.outgoing.flow-events-out.*` + (connector=smallrye-rabbitmq、exchange.type=topic/durable)、`enabled` 绑定同一开关、 + `quarkus.rabbitmq.devservices.enabled=false`。 +- 埋点(11 处):DefinitionProcessor (create/update/deploy/delete)、FlowExecutor (preExecute→started、execute finally + FAILED→failed、postExecute/postCommit COMPLETED/END→completed)、RuntimeProcessor.terminateProcess + (TERMINATED→terminated,子流程级联逐个)、UserTaskExecutor + (doExecute→suspended、postCommit→committed、doRollback→rollback.suspended)。 + +**测试:** + +- `FlowEventPublisherTest`(2):本地 EventBus 广播 + 字段完整 + null no-op。 +- `FlowEventRabbitmqBridgeTest`(2):默认禁用不连 broker(应用无 broker 启动即证)、禁用桥接对本地事件透明。未采用 SmallRye + InMemoryConnector(当前 Quarkus 3.33.1 未提供配套 in-memory 扩展,InMemorySinkImpl 无 bean 定义注解,注入不满足)。 + +**验证:** + +- `mvn clean compile -pl '!flexmodel-engine/flexmodel-maven-plugin'` → BUILD SUCCESS +- +`mvn test -pl flexmodel-server -Dtest=FlowEventPublisherTest,FlowEventRabbitmqBridgeTest,DefinitionProcessorTest,RuntimeProcessorTest` → +25 tests, 0 failures, 0 errors(含 DefinitionProcessorTest 4、RuntimeProcessorTest 17 回归通过) + +**备注:** RabbitMQ 转发(enabled=true)端到端需真实 broker 或 Testcontainers,按计划保持可选、默认关闭;v1 不含 +node-instance-created/service-task 等事件,后续同模式按需追加。 + +## 去掉 flexmodel.flow.events.rabbitmq.enabled 业务开关(2026-08-19) + +**背景:** 该业务开关与 SmallRye 通道 `mp.messaging.outgoing.flow-events-out.enabled` 重复;且桥接 `enabled` 默认 false +时即便通道启用也不转发,属冗余控制层。 + +**变更:** + +- `application.properties`:删除 `flexmodel.flow.events.rabbitmq.enabled`,通道 `enabled=false` 成为单一控制(启用置 true + + broker 连接配置)。 +- `FlowEventRabbitmqBridge`:删除 `@ConfigProperty enabled` 字段、11 处 `if(!enabled) return;` early-return、`isEnabled()` + 测试探针;转发完全由 `Instance` 解析性决定——通道禁用时 SmallRye 注入 no-op emitter,`forward()` 发送即丢弃、不连 + broker。 +- 删除 `FlowEventConfig`(`@ConfigMapping` 不再需要声明已移除的配置根)。 +- `FlowEventRabbitmqBridgeTest`:移除 `isEnabled/isEmitterResolvable` 断言,改为断言桥接 bean 存在且默认通道禁用时发布不抛出、本地事件照常广播。 + +**验证:** + +- `mvn test-compile -pl '!flexmodel-engine/flexmodel-maven-plugin'` → ExitCode 0 +- +`mvn test -pl flexmodel-server -Dtest=FlowEventPublisherTest,FlowEventRabbitmqBridgeTest,DefinitionProcessorTest,RuntimeProcessorTest` → +25 tests, 0 failures, 0 errors + +**设计效果:** 单一配置源——`mp.messaging.outgoing.flow-events-out.enabled` 同时控制是否连 broker 与是否转发;默认 false +零侵入,本地 EventBus 事件始终发布。 + +## Review 修复:桥接静默跳过 + variables 防御性拷贝(2026-08-19) + +**P1 桥接默认配置逐事件 WARN 栈:** `Instance.isUnsatisfied()` 对禁用通道返回 false,`.get()` 抛 SRMSG00019 被捕获记 +WARN(含事件载荷/流程变量)。修复:`FlowEventRabbitmqBridge` 注入 +`@ConfigProperty("mp.messaging.outgoing.flow-events-out.enabled", defaultValue="false") channelEnabled`,`forward()` 首行 +`if(!channelEnabled) return;` 静默跳过;转发失败日志改为仅记 routingKey + payloadType,不再打印整个事件(避免泄露流程变量)。 + +**P2 发布可变 variables 快照并发风险:** `FlowInstanceStartedEvent`/`FlowInstanceCompletedEvent`/`UserTaskSuspendedEvent` +构造时按引用持有 `runtimeContext.getInstanceDataMap()`,发布线程后续 mutate 同一 map,异步消费者/序列化读取共享 map 可能不一致或 +CME。修复:三个事件构造器对 variables 做 `new HashMap<>(variables)` 防御性拷贝(null 安全),发布快照不可变。 + +**验证:** `mvn test-compile` 通过; +`FlowEventPublisherTest(2)/FlowEventRabbitmqBridgeTest(2)/DefinitionProcessorTest(4)/RuntimeProcessorTest(17)` 共 25 测试 +0 失败 0 错误。默认配置下桥接不再逐事件打 WARN(SRMSG00232 通道禁用为启动一次性诊断)。 + +## Flow 生命周期事件 RabbitMQ Testcontainers E2E 测试(2026-08-19) + +**目标:** 增加 Testcontainers 端到端测试,启动真实 RabbitMQ broker 验证 `FlowEventRabbitmqBridge` 以正确 routing key + +JSON 载荷推送事件到 topic 交换机。 + +**新增/修改:** + +- `flexmodel-server/pom.xml`:加 `org.testcontainers:rabbitmq:${testcontainers.version}`(test scope,版本由父 pom + `testcontainers.version=1.21.4` 管理;amqp-client 5.x 传递可用)。 +- `RabbitMqTestResource.java`(前序会话已写):实现 `QuarkusTestResourceLifecycleManager`,启动 + `RabbitMQContainer("rabbitmq:3-management")`,注入 `quarkus.rabbitmq.host/port/username/password` 与 + `mp.messaging.outgoing.flow-events-out.enabled=true`,暴露 static host/port/username/password 供测试取连接坐标。 +- `FlowEventRabbitmqBridgeE2ETest.java`(新增):`@QuarkusTest` + `@QuarkusTestResource(SQLiteTestResource.class)` + + `@QuarkusTestResource(value=RabbitMqTestResource.class, restrictToAnnotatedClass=true)`。两个用例: + - `flowDeployedEventForwardedToBroker`:amqp-client 临时队列绑定交换机 `flexmodel.flow.events`(routing key + `flow.deployed`),发布 `FlowDeployedEvent`,`basicGet` 轮询(≤15s)断言 routing key 与 JSON + 字段(projectId/caller/flowModuleId/flowDeployId/timestamp)。 + - `flowInstanceStartedEventWithVariablesForwarded`:同模式断言 `flow.instance.started` 与 variables + 快照(amount/approved)经 JSON 序列化完整。 + +**关键设计决策:** + +- **资源泄漏修复(重要):** 初版 `@QuarkusTestResource(RabbitMqTestResource.class)` 默认 `restrictToAnnotatedClass=false` + ,导致 broker 资源被 Quarkus 全局应用到所有共享应用上下文的测试——即便未选 E2E 测试,其他测试(`FlowEventPublisherTest` + 等)也触发 `rabbitmq:3-management` 镜像拉取,无 Docker Hub 环境下整套测试失败。改为 `restrictToAnnotatedClass=true` + ,资源仅对本测试类生效。验证确认:非 E2E 测试不再拉取镜像。 +- **opt-in 开关:** E2E 测试加 `@EnabledIfEnvironmentVariable(named="FLEXMODEL_E2E_RABBITMQ", matches="true")` + ,默认跳过。项目此前无任何 Docker 依赖测试,此为首个;默认 `mvn test` 在无 Docker/无 Docker Hub 环境保持绿色。运行需:本机 + Docker 可用 + 能拉取 `rabbitmq:3-management` + 设环境变量 `FLEXMODEL_E2E_RABBITMQ=true`。 +- routing key 取自 AMQP envelope(`routingKey()` 是方法不进 JSON),JSON 载荷用 Jackson 解析;带 variables + 的事件经防御性拷贝保证快照不可变(见前序 Review 修复)。 + +**验证:** + +- `mvn test-compile -pl '!flexmodel-engine/flexmodel-maven-plugin' -q` → ExitCode 0。 +- 综合回归 + `mvn test -pl flexmodel-server -Dtest=FlowEventPublisherTest,FlowEventRabbitmqBridgeTest,FlowEventRabbitmqBridgeE2ETest,DefinitionProcessorTest,RuntimeProcessorTest` → + Tests run: 27, Failures: 0, Errors: 0, Skipped: 2(E2E 默认跳过),BUILD SUCCESS,无 Docker 镜像拉取。 +- **E2E 实跑未完成:** 当前环境 Docker Hub 不可达(`registry-1.docker.io` EOF / 配置镜像源 `docker.1panel.live` 超时), + `rabbitmq:3-management` 无法拉取,故 `FLEXMODEL_E2E_RABBITMQ=true` 实跑未通过。代码逻辑正确、编译通过;待网络恢复或换可用镜像源后,设该环境变量即可执行端到端验证。 + +**未决/后续:** + +- E2E 实跑待 Docker Hub 可达后补验证(设 `FLEXMODEL_E2E_RABBITMQ=true` 运行)。 +- 若 CI 需常态化跑 E2E,建议在 CI 配置可用 RabbitMQ 镜像源或预拉镜像,并设该环境变量。 + +## E2E 测试修复:RabbitMQ 连接配置键(2026-08-19 续) + +**问题:** 首次运行 `FlowEventRabbitmqBridgeE2ETest` 失败——应用侧 SmallRye outgoing channel `Connection refused`,事件未发到 +broker,测试轮询超时。日志报 `Unrecognized configuration key "quarkus.rabbitmq.password"`。 + +**根因:** 反编译 `quarkus-messaging-rabbitmq` 扩展的 Quarkus config root `RabbitMQBuildTimeConfig` 确认:其仅注册 +`devservices`、`credentialsProvider`、`credentialsProviderName` 字段, **不注册 `host/port/username/password` 连接字段**。故 +`RabbitMqTestResource` 注入的 `quarkus.rabbitmq.host/port/username/password` 全部无效(unrecognized),SmallRye client 回退默认 +`localhost:5672`,连不上 Testcontainers 随机映射端口。这些连接字段由 SmallRye connector 自身读取(channel 级 +`mp.messaging.outgoing..host/port/username/password`,或全局别名 `rabbitmq-host` 等),通过查 +`smallrye-reactive-messaging-rabbitmq-4.33.0` 源码 `RabbitMQConnectorCommonConfiguration` 确认。 + +**修复:** `RabbitMqTestResource.start()` 改用 channel 级连接配置: +`mp.messaging.outgoing.flow-events-out.host/port/username/password`;保留 `quarkus.rabbitmq.devservices.enabled=false`(防 +DevServices 自启)与 `mp.messaging.outgoing.flow-events-out.enabled=true`。测试侧 amqp-client 仍用 static `host/port` +订阅交换机。 + +**验证:** 设 `FLEXMODEL_E2E_RABBITMQ=true` 运行 → +`SRMSG17036: RabbitMQ broker configured to [localhost:52374] for channel flow-events-out` + +`SRMSG17007: Connection with RabbitMQ broker established` → **Tests run: 2, Failures: 0, Errors: 0, Skipped: 0, BUILD +SUCCESS**。两个用例(`flowDeployedEventForwardedToBroker`、`flowInstanceStartedEventWithVariablesForwarded`)断言 routing +key 与 JSON 载荷字段(含 variables 快照)完整通过。 + +**结论:** E2E 端到端验证完成。Flow 生命周期事件经 `FlowEventPublisher`→EventBus→`FlowEventRabbitmqBridge`→SmallRye +outgoing channel→RabbitMQ topic 交换机链路,routing key 与 JSON 载荷均正确。默认 `mvn test` 不依赖 Docker(E2E opt-in 跳过), +`restrictToAnnotatedClass=true` 确保 broker 资源不泄漏到其他测试。 + +## UserTask 事件携带 nodeAttributes + 文档补全数据结构(2026-08-19) + +**需求:** 订阅者常需解析节点定义里的扩展属性做业务(审批人、表单、阈值等)。外部 RabbitMQ 订阅者访问不到定义仓库,故将节点属性快照随事件携带。 + +**改动:** + +- 事件类:`UserTaskSuspendedEvent`/`UserTaskCommittedEvent`/`UserTaskRollbackSuspendedEvent` 各加 `Map nodeAttributes` 字段,构造器做 `new HashMap<>(nodeAttributes)` 防御性拷贝(与 variables 一致)。 +- `UserTaskExecutor` 三处埋点填充 `nodeAttributes`:`doExecute` 用 `flowElement.getProperties()`;`postCommit` 用 `runtimeContext.getCurrentNodeModel().getProperties()`;`doRollback` 用 `FlowModelUtil.getFlowElement(flowElementMap, nodeKey).getProperties()`。 +- 测试:`FlowEventRabbitmqBridgeE2ETest` 新增 `userTaskSuspendedEventWithNodeAttributesForwarded`,验证 routing key + variables + nodeAttributes 经 broker 转发后 JSON 载荷完整(含中文字段名、boolean)。 +- 文档 `flexmodel-website/docs/tutorial/features/flow.md`: + - 修正公共字段表——`routingKey` 不在 JSON 载荷(是方法非字段),改注其随 AMQP envelope 投递;同步修正 Python 订阅示例从 `method.routing_key` 读取。 + - 事件总览表三个 UserTask 事件加 `nodeAttributes`。 + - 新增「事件数据结构」小节,按定义层/实例层/用户任务层分组给出每种事件的 JSON 载荷示例。 + +**验证:** + +- `mvn test-compile` → ExitCode 0。 +- E2E(`FLEXMODEL_E2E_RABBITMQ=true`)→ Tests run: 3, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS(含新增 nodeAttributes 用例)。 +- 默认回归 → Tests run: 28, Failures: 0, Errors: 0, Skipped: 3(E2E 默认跳过),BUILD SUCCESS。 + +**设计要点:** nodeAttributes = 节点定义的 `FlowElement.properties` 快照,事件发生时已确定且不可变;外部订阅者无需回查定义仓库即可读取节点配置,解耦核心与外部系统。