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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 5 additions & 5 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -105,12 +105,12 @@ WINDUP_MQ_MAX_CONSUME_ATTEMPTS=5
WINDUP_MQ_CONSUME_LEASE_SECONDS=1800
WINDUP_MQ_EMAIL_HANDLER_RETRIES=3
# worker 内 handler 并行度(单进程内 ThreadPoolExecutor)
# 生成期间不占 Postgres 连接,可按机器内存与上游配额上调
# 图 / 动作分两条 Stream 独立出队。默认对齐现网,避免无环境变量时打满上游
WINDUP_MQ_EMAIL_CONCURRENCY=8
WINDUP_MQ_GENERATION_IMAGE_CONCURRENCY=16
WINDUP_MQ_GENERATION_ACTION_CONCURRENCY=8
# poll 使用独立线程池,不与图/动作共享槽位
WINDUP_MQ_GENERATION_POLL_CONCURRENCY=16
WINDUP_MQ_GENERATION_IMAGE_CONCURRENCY=4
WINDUP_MQ_GENERATION_ACTION_CONCURRENCY=2
# poll 使用独立线程池,不与图/动作提交共享槽位
WINDUP_MQ_GENERATION_POLL_CONCURRENCY=2
# i2v 轮询走 ZSET 延迟队列,到期促进到 Stream;不用 Redis 过期事件
WINDUP_MQ_DELAYED_ZSET=windup:zset:delayed
WINDUP_MQ_DELAYED_CLAIM_LIMIT=50
Expand Down
33 changes: 24 additions & 9 deletions backend/packages/app/src/windup_app/bootstrap/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,13 @@
import threading
import time

from windup_app.server.mq.catalog import all_stream_specs, email_stream_spec, generation_stream_spec
from windup_app.server.mq.catalog import (
all_stream_specs,
email_stream_spec,
generation_action_stream_spec,
generation_image_stream_spec,
generation_stream_spec,
)
from windup_app.server.orchestrator import task_repo
from windup_app.server.orchestrator.executor import (
bind_matte,
Expand Down Expand Up @@ -67,25 +73,31 @@ def _handle_signal(_signum, _frame) -> None:
signal.signal(signal.SIGTERM, _handle_signal)
signal.signal(signal.SIGINT, _handle_signal)

consumers = [
StreamConsumer(
email_stream_spec(),
def _generation_consumer(spec):
return StreamConsumer(
spec,
run_image_task=run_image_task,
run_action_task=run_action_task,
run_direction_set_task=run_direction_set_task,
run_view_sheet_task=run_view_sheet_task,
stop_event=stop_event,
),
resume_action_poll=resume_action_poll,
resume_action_client_bake=resume_action_client_bake,
)

consumers = [
StreamConsumer(
generation_stream_spec(),
email_stream_spec(),
run_image_task=run_image_task,
run_action_task=run_action_task,
run_direction_set_task=run_direction_set_task,
run_view_sheet_task=run_view_sheet_task,
stop_event=stop_event,
resume_action_poll=resume_action_poll,
resume_action_client_bake=resume_action_client_bake,
),
_generation_consumer(generation_image_stream_spec()),
_generation_consumer(generation_action_stream_spec()),
# 过渡 drain:切流前已进旧 generation Stream 的消息还要被消费。
_generation_consumer(generation_stream_spec()),
]
threads = [consumer.start() for consumer in consumers]
relay_thread = start_relay_loop(stop_event)
Expand All @@ -99,7 +111,10 @@ def _handle_signal(_signum, _frame) -> None:
)
pending_thread.start()

logger.info("windup worker 已启动 | streams=%s", [s.stream for s in all_stream_specs()])
logger.info(
"windup worker 已启动 | streams=%s",
[s.stream for s in (*all_stream_specs(), generation_stream_spec())],
)

try:
while not stop_event.is_set():
Expand Down
85 changes: 68 additions & 17 deletions backend/packages/app/src/windup_app/server/mq/catalog.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
- ``recover_as``:PENDING 任务按此 GenerationType 值重入队;轮询类留 None
3. 在 handlers 的 ``HANDLERS`` 登记可调用对象

图片与动作分两条 Stream,出队互不堵。旧 ``GENERATION_STREAM`` 只给过渡 drain。
不在此表里塞积分账本或 SSE EventBus。是否新开 Stream 仍按 SLA 决定。
"""

Expand All @@ -31,9 +32,13 @@

EMAIL_STREAM = "windup:stream:email"
GENERATION_STREAM = "windup:stream:generation"
GENERATION_IMAGE_STREAM = "windup:stream:generation-image"
GENERATION_ACTION_STREAM = "windup:stream:generation-action"

EMAIL_GROUP = "email"
GENERATION_GROUP = "generation"
GENERATION_IMAGE_GROUP = "generation-image"
GENERATION_ACTION_GROUP = "generation-action"

MSG_TYPE_VERIFICATION_CODE = "verification_code"
MSG_TYPE_CHARACTER_IMAGE = "character_image"
Expand Down Expand Up @@ -72,17 +77,17 @@ class TypeSpec:


def generation_image_concurrency() -> int:
# executor 已短 session:生成/上传不再占连接,默认不再被 15 连接池卡住
return _env_int("WINDUP_MQ_GENERATION_IMAGE_CONCURRENCY", 16)
# 默认对齐现网:图与动作分队后各自限流,避免无环境变量时按 16/8 打上游
return _env_int("WINDUP_MQ_GENERATION_IMAGE_CONCURRENCY", 4)


def generation_action_concurrency() -> int:
return _env_int("WINDUP_MQ_GENERATION_ACTION_CONCURRENCY", 8)
return _env_int("WINDUP_MQ_GENERATION_ACTION_CONCURRENCY", 2)


def generation_poll_concurrency() -> int:
# 单次 inspect + 偶尔下载,不 sleep,默认高于 action 建单并发
return _env_int("WINDUP_MQ_GENERATION_POLL_CONCURRENCY", 16)
# 单次 inspect + 偶尔下载,不 sleep;与动作提交分池,不占提交槽
return _env_int("WINDUP_MQ_GENERATION_POLL_CONCURRENCY", 2)


def type_specs() -> tuple[TypeSpec, ...]:
Expand All @@ -96,30 +101,30 @@ def type_specs() -> tuple[TypeSpec, ...]:
),
TypeSpec(
msg_type=MSG_TYPE_CHARACTER_IMAGE,
stream=GENERATION_STREAM,
stream=GENERATION_IMAGE_STREAM,
pool=POOL_SHARED,
concurrency=generation_image_concurrency(),
limit=True,
recover_as=MSG_TYPE_CHARACTER_IMAGE,
),
TypeSpec(
msg_type=MSG_TYPE_CHARACTER_ACTION,
stream=GENERATION_STREAM,
stream=GENERATION_ACTION_STREAM,
pool=POOL_SHARED,
concurrency=generation_action_concurrency(),
limit=True,
recover_as=MSG_TYPE_CHARACTER_ACTION,
),
TypeSpec(
msg_type=MSG_TYPE_CHARACTER_ACTION_POLL,
stream=GENERATION_STREAM,
stream=GENERATION_ACTION_STREAM,
pool=POOL_POLL,
concurrency=generation_poll_concurrency(),
limit=True,
),
TypeSpec(
msg_type=MSG_TYPE_CHARACTER_ACTION_CLIENT_BAKE,
stream=GENERATION_STREAM,
stream=GENERATION_ACTION_STREAM,
pool=POOL_POLL,
concurrency=generation_poll_concurrency(),
limit=True,
Expand All @@ -135,7 +140,23 @@ def type_spec(msg_type: str) -> TypeSpec | None:


def types_for_stream(stream: str) -> tuple[TypeSpec, ...]:
return tuple(spec for spec in type_specs() if spec.stream == stream)
matched = tuple(spec for spec in type_specs() if spec.stream == stream)
if stream != GENERATION_STREAM:
return matched
# 旧产线 Stream 过渡 drain:切流后仍要按图/动作类型处理残留消息。
drained = tuple(
spec
for spec in type_specs()
if spec.stream in (GENERATION_IMAGE_STREAM, GENERATION_ACTION_STREAM)
)
return matched + drained


def stream_for_msg_type(msg_type: str) -> str:
spec = type_spec(msg_type)
if spec is None:
raise ValueError(f"未知消息类型: {msg_type}")
return spec.stream


def msg_type_for_generation(task_type: str) -> str:
Expand All @@ -157,8 +178,8 @@ def msg_type_for_generation(task_type: str) -> str:
def _pool_size(stream: str, pool: str) -> int:
return sum(
spec.concurrency
for spec in type_specs()
if spec.stream == stream and spec.pool == pool
for spec in types_for_stream(stream)
if spec.pool == pool
)


Expand All @@ -170,22 +191,44 @@ def email_stream_spec() -> StreamSpec:
)


def generation_image_stream_spec() -> StreamSpec:
return StreamSpec(
stream=GENERATION_IMAGE_STREAM,
group=GENERATION_IMAGE_GROUP,
concurrency=_pool_size(GENERATION_IMAGE_STREAM, POOL_SHARED),
)


def generation_action_stream_spec() -> StreamSpec:
return StreamSpec(
stream=GENERATION_ACTION_STREAM,
group=GENERATION_ACTION_GROUP,
concurrency=_pool_size(GENERATION_ACTION_STREAM, POOL_SHARED),
)


def generation_worker_pool_size() -> int:
# image/action 共用一个线程池。poll 走独立 pool 名,不能加进这个数字,
# 否则 image 占满线程后 poll 只能排队。
return _pool_size(GENERATION_STREAM, POOL_SHARED)
# 旧 Stream drain 的共享池:图+动作合计。poll 不进这个数字。
return _pool_size(GENERATION_IMAGE_STREAM, POOL_SHARED) + _pool_size(
GENERATION_ACTION_STREAM, POOL_SHARED
)


def generation_stream_spec() -> StreamSpec:
"""旧 ``generation`` Stream 的过渡 drain 规格。新任务不要再往这里投。"""
return StreamSpec(
stream=GENERATION_STREAM,
group=GENERATION_GROUP,
concurrency=generation_worker_pool_size(),
)


def all_stream_specs() -> tuple[StreamSpec, StreamSpec]:
return email_stream_spec(), generation_stream_spec()
def all_stream_specs() -> tuple[StreamSpec, ...]:
return (
email_stream_spec(),
generation_image_stream_spec(),
generation_action_stream_spec(),
)


__all__ = [
Expand All @@ -194,9 +237,14 @@ def all_stream_specs() -> tuple[StreamSpec, StreamSpec]:
"EMAIL_HANDLER_RETRIES",
"GENERATION_STREAM",
"GENERATION_GROUP",
"GENERATION_IMAGE_STREAM",
"GENERATION_IMAGE_GROUP",
"GENERATION_ACTION_STREAM",
"GENERATION_ACTION_GROUP",
"GENERATION_PENDING_MAX_AGE_SECONDS",
"GENERATION_RUNNING_STALE_SECONDS",
"MSG_TYPE_CHARACTER_ACTION",
"MSG_TYPE_CHARACTER_ACTION_CLIENT_BAKE",
"MSG_TYPE_CHARACTER_ACTION_POLL",
"MSG_TYPE_CHARACTER_IMAGE",
"MSG_TYPE_VERIFICATION_CODE",
Expand All @@ -208,11 +256,14 @@ def all_stream_specs() -> tuple[StreamSpec, StreamSpec]:
"all_stream_specs",
"email_stream_spec",
"generation_action_concurrency",
"generation_action_stream_spec",
"generation_image_concurrency",
"generation_image_stream_spec",
"generation_poll_concurrency",
"generation_stream_spec",
"generation_worker_pool_size",
"msg_type_for_generation",
"stream_for_msg_type",
"type_spec",
"type_specs",
"types_for_stream",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,8 @@
from dataclasses import asdict, dataclass

from windup_app.server.mq.catalog import (
GENERATION_STREAM,
MSG_TYPE_CHARACTER_ACTION_CLIENT_BAKE,
stream_for_msg_type,
)
from windup_framework.db.redis import get_redis
from windup_framework.mq.delayed import schedule_delayed
Expand Down Expand Up @@ -93,7 +93,7 @@ def open_job(task_id: int, spec: ClientBakeSpec) -> float:
pipe.execute()
schedule_delayed(
delay_s=DEADLINE_S,
stream=GENERATION_STREAM,
stream=stream_for_msg_type(MSG_TYPE_CHARACTER_ACTION_CLIENT_BAKE),
msg_type=MSG_TYPE_CHARACTER_ACTION_CLIENT_BAKE,
payload={"task_id": task_id, "reason": REASON_TIMEOUT},
dedupe_key=f"generation:{task_id}:clientbake:timeout",
Expand Down Expand Up @@ -169,7 +169,7 @@ def schedule_resume(task_id: int, reason: str = REASON_FRAMES, detail: str = "")
payload["detail"] = detail[:200]
schedule_delayed(
delay_s=0,
stream=GENERATION_STREAM,
stream=stream_for_msg_type(MSG_TYPE_CHARACTER_ACTION_CLIENT_BAKE),
msg_type=MSG_TYPE_CHARACTER_ACTION_CLIENT_BAKE,
payload=payload,
dedupe_key=f"generation:{task_id}:clientbake:{reason}",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
from dataclasses import dataclass
from typing import Any

from windup_app.server.mq.catalog import GENERATION_STREAM, MSG_TYPE_CHARACTER_ACTION_POLL
from windup_app.server.mq.catalog import MSG_TYPE_CHARACTER_ACTION_POLL, stream_for_msg_type
from windup_app.server.mq.i2v_state import (
I2V_FIRST_POLL_S,
I2V_MAX_WAIT_S,
Expand Down Expand Up @@ -91,7 +91,7 @@ def schedule(
)
schedule_delayed(
delay_s=wait,
stream=GENERATION_STREAM,
stream=stream_for_msg_type(MSG_TYPE_CHARACTER_ACTION_POLL),
msg_type=MSG_TYPE_CHARACTER_ACTION_POLL,
payload=_poll_payload(task_id, poll_count),
dedupe_key=_poll_dedupe(task_id, poll_count),
Expand Down Expand Up @@ -149,7 +149,7 @@ def reschedule_if_waiting(task_id: int, *, delay_s: float = 1) -> bool:
poll_count = int(state.get("poll_count") or 0)
schedule_delayed(
delay_s=delay_s,
stream=GENERATION_STREAM,
stream=stream_for_msg_type(MSG_TYPE_CHARACTER_ACTION_POLL),
msg_type=MSG_TYPE_CHARACTER_ACTION_POLL,
payload=_poll_payload(task_id, poll_count),
dedupe_key=_poll_dedupe(task_id, poll_count),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,8 @@

from windup_app.server.mq.catalog import (
GENERATION_RUNNING_STALE_SECONDS,
GENERATION_STREAM,
msg_type_for_generation,
stream_for_msg_type,
)
from windup_app.server.orchestrator import billing, task_repo
from windup_app.server.orchestrator.i2v_poll import reschedule_if_waiting
Expand Down Expand Up @@ -109,10 +109,11 @@ def _requeue_pending(
if hasattr(task.task_type, "value")
else str(task.task_type)
)
msg_type = msg_type_for_generation(task_type)
message_id = publisher.enqueue(
session,
stream=GENERATION_STREAM,
msg_type=msg_type_for_generation(task_type),
stream=stream_for_msg_type(msg_type),
msg_type=msg_type,
payload={
"task_id": task.id,
"task_type": task_type,
Expand Down
4 changes: 2 additions & 2 deletions backend/packages/app/src/windup_app/web/api/generation.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,8 @@
from windup_app.server.character.model import Character, CharacterData
from windup_app.server.orchestrator import billing, task_repo
from windup_app.server.mq.catalog import (
GENERATION_STREAM,
msg_type_for_generation,
stream_for_msg_type,
)
from windup_app.server.orchestrator.service import service as generation_service
from windup_app.server.orchestrator.model import (
Expand Down Expand Up @@ -531,7 +531,7 @@ def _publish_generation_after_commit(
msg_type = msg_type_for_generation(task_type)
message_id = publisher.enqueue(
session,
stream=GENERATION_STREAM,
stream=stream_for_msg_type(msg_type),
msg_type=msg_type,
payload={"task_id": task_id, "task_type": task_type},
dedupe_key=dedupe_key or f"generation:{task_id}",
Expand Down
Loading
Loading