makepad/libs/ai/hub/python/music3_worker.py
Admin e37c263b9c ai backbone: the hub era — makepad_ai deleted, every backend is a hub pipe; machine residency elections, job leases, ETA placement; the store only stores; creator pipelines run in the app (aicore)
Squashed from work:
- asset-ai: FastH3 4-step fast video backend; clip keyframes on the wire
- asset-ui: loop video chains — text→image→video that ends where it began
- h3: safetensors -> pruned-Q4_K GGUF quantizer for the 24GB DiT tiers
- h3_quant_gguf verify: row-error gates calibrated to the measured Q4_K floor
- asset-ai realtime: the feedback loop — the source anchors, the drifted frame inits
- asset-ai realtime: a feedback loop survives a resize and travels by default
- asset-ai realtime: the feedback loop frees itself from the feed handshake and pauses for its listener
- asset-ai realtime: the outbound encode leaves the loop's critical path
- asset-ai ocr: the ocr domain — Chandra 2 at page resolution, and the tower goes planner-owned
- llm slots: a lane can hold an image span — embedding prefill and a rope cursor of its own
- vision tower on CUDA: the encode leg gets its two missing kernels
- llm/ocr: one M-RoPE grid encoder for both image paths, and a livelock made an error
- vision tower on CUDA: the f16 GEMM keeps the precision it was throwing away
- live: a feed that moves box takes its trip with it — one seed image
- vision tower on CUDA: the tiled attention becomes bit-exact, and tensor cores go
- llm prefill on CUDA: the MMA attention kernel gets the tile a 4-to-1 model needs
- asset-ai ocr: the CUDA encode lane joins the integration — vision-parity sits beside run's three arms, and the kernels
- Merge branch 'ocr-perf-integration' into work
- asset-ai: the live anchor can follow the trip, and text leaves the 5090
- asset-ai: the camera moves the world, and the world starts still
- asset-import: the EA strategy classics, in the one 2D contract
- rtsmap: one seeded generator for tiled strategy maps
- asset-ui: one card for the strategy classics, with a pack dropdown
- asset-ai: music3 reference-audio path, ocr/h3 backends, registry
- asset: mp4 sample index for range-streaming, chat tools, import profiles
- cnc: tiberium is twelve growth frames, not twelve empty variants
- platform: native file and save dialogs, in-house on all three desktops
- chat: the scan holds out for a lane home
- chat: a full home queues you — take the free lane
- chat: the preload has a percentage, and the boundless cap stops showing
- llm cuda: the 32x2 attention tile — even GQA ratios stay on MMA
- sa3 gets a bake path: the sfx model's tables precomputed by a diffusion-side bin
- sqlite_query: anti-join regression test
- td import: HARV's second frame block is its harvesting cycle, not a turret
- asset-ui: sprite enhancement runs on the 32B dev DiT — distillation, not the prompt, was the ceiling
- ai-hub: makepad-asset-ai becomes makepad-ai-hub at libs/ai/hub, the chat pane becomes makepad-chat-ui, the service bin
- asset-ui: test health fixtures grow the realtime field they were born without
- ai-hub: one home at ~/.makepad — weights/ run/ cache/ logs/, the service cache migrates from ai_content by a single re
- ai-hub: subprocess workers die with the node — process groups everywhere, PDEATHSIG on linux, one KILL_ON_JOB_CLOSE Jo
- ai-hub: the hub object — AiHub::in_process, pipes vocabulary, and the local LLM engine generalized out of mpfiles (aic
- strict-json: the dependency-free JSON module gets its own crate; asset-client re-exports it so nothing downstream move
- ai-hub: the machine layer — node entries, the 0600 machine token, and the residency election that IS the lock (aicore
- ai-hub: MPHUB1 — the fabric beacon only dedicated nodes can send (aicore §4)
- ai-hub: job leases — work lives only while it is renewed (aicore §8)
- asset-creator: the pipeline library is born — specs, the deps gate, and the derived-state law (aicore §9)
- ai-hub: RAM residency facts — the CPU-side twin of residency.rs (aicore §3)
- ai-hub: ETA placement primitives — relative GPU throughput, the four-term estimate, and an observable breakdown (aicor
- ai-hub: leases go live on the wire — origin fields on submit, /job/<id>/keepalive, /bye, and the reaper that cancels w
- ai-hub: the chat providers move in — fleet qwen, openai, grok, claude/codex/grok CLIs, the responses driver, and the w
- asset-creator: the engine — one pipeline run against the hub, deps-gated, spliced, cancellable, resumable-by-construct
- ai-hub: the machine node mode — --machine binds loopback, registers in ~/.makepad/run, and exits on its own once idle
- asset-creator: makepad-creator-run — the detached client for runs that must outlive a window (aicore §9)
- ai-hub: a native Claude Messages-API provider — API-key or Claude Code OAuth, bounded SSE streaming, injected tools (a
- route + converse: off makepad_ai — the Agent seam moves to converse, route's cloud dispatcher rides the hub's Claude p
- asset-creator: the preset tables move in — fifteen chain-policy constants shared by every creator app (aicore §9 / P6)
- makepad_ai is deleted — every backend is a hub pipe, the agent seam lives with its consumers (aicore §14, decided 2026
- ai-hub: loads hold the machine residency election — set_model_state claims on Loaded and publishes the service port (a
- ai-hub: chats run the machine election — route to a serving holder, wait on a loading one, claim and publish when open
- ai-hub: pick_for_domain_eta — ETA-ranked placement over the shared hard-filter core (aicore §6 / P4)
- asset-creator: the engine picks a provider per stage at dispatch time — a chain's later stages see fresh fleet state (
- ai-hub: the fabric secret gates the service HTTP surface — bearer on everything but /health and the ticketed peer path
- vj: DREAM runs execute in the app — pipelines.rs becomes the run it used to watch (aicore §9 / F1)
- asset-creator: the runner — generate one thing and put it in the catalog, one implementation for every surface (aicore
- chat-ui: the session runs in the app — no broker anywhere on the chat path (aicore P8 / F5)
- asset-store: assets.query is a first-class query endpoint — the bounded SQL surface outlives the broker (aicore P8 / F
- asset-creator: CreatorTools — the chat tool pack for a store that only stores (aicore §9 / P8)
- asset-store: the shrink — the store stores (aicore P7)
- importer + asset-server host: the coordination era ends (aicore P7)
- store config purge + asset-ui goes fleet-direct; the derive protocol gets its route proof (aicore P7)
- client + chat dispatcher: the dead wire comes out (aicore P7/P8)
- ai-hub: 0.3.0 — the health version says which era a node runs
- ai-hub: the default fleet is 'gen' — apps hear the LAN without env plumbing
- ai-hub: the preload note percents the prefill, not the job bar
- ai-hub: conversations keep their KV — the wire mirror, the lane identity, the in-turn dynamic context (aicore §7)
- ai-hub: an open-think model is thinking from its first token
- libs: the zero-warning sweep — stitch casts say what they mean, xatlas keeps upstream's surface quietly
- zero-warning sweep, round two — the first full-workspace pass
- zero-warning sweep, round three — the model lanes and the deep examples
- zero-warning sweep, round four — the last stragglers
- zero-warning sweep, round five — vj and chat-ui
- zero-warning sweep, round six — three cascades

Co-authored-by: Claude <info@makepad.nl>
2026-09-01 16:46:31 +02:00

280 lines
11 KiB
Python

# music3_worker.py - persistent MiniMax-Music3 job worker for the
# makepad-ai-content `music` domain backend (music3_backend.rs). Line protocol
# (same shape as fw_worker.py):
# stdin : one JSON object per line: {"prompt": str, "lyrics": str,
# "duration_s": float, "seed": int, "out_wav": str} or {"exit": true}
# stdout: events prefixed "@EV " (everything else is ignored by the parent):
# {"ev":"stage","stage":name[,"k":i,"n":total]}
# {"ev":"ready"} after model load
# {"ev":"done","wav":path} job finished
# {"ev":"error","message":text} load/job failed (worker lives
# on after job errors)
#
# Runtime: the official diffusers ModularPipeline integration
# (MiniMaxAI/MiniMax-Music3 modular_model_index.json; diffusers PR #14456,
# pinned venv commit dafe3733fcfdbf3c48915fe77be3aef65b5d6a2d per the model
# card). Weights come from the service cache (registry-managed download), so
# the worker runs fully hub-offline against --model-dir.
#
# The recorded component specs in modular_model_index.json point at the hub
# repo id ("MiniMaxAI/MiniMax-Music3"); to guarantee offline loads we build a
# hardlinked view of the model dir with those specs rewritten to the local
# view path, then from_pretrained() the view.
import argparse
import json
import os
import shutil
import sys
os.environ.setdefault("HF_HUB_OFFLINE", "1")
os.environ.setdefault("TRANSFORMERS_OFFLINE", "1")
os.environ.setdefault("HF_HUB_DISABLE_TELEMETRY", "1")
def ev(**kw):
sys.stdout.write("@EV " + json.dumps(kw) + "\n")
sys.stdout.flush()
def build_local_view(model_dir, view_dir):
"""Hardlink (fallback copy) the cached model into view_dir, rewriting the
modular_model_index.json component sources to the view path so every
component loads from local disk regardless of how diffusers resolves the
recorded hub repo id."""
for root, _dirs, files in os.walk(model_dir):
rel = os.path.relpath(root, model_dir)
dst_root = os.path.join(view_dir, rel) if rel != "." else view_dir
os.makedirs(dst_root, exist_ok=True)
for name in files:
src = os.path.join(root, name)
dst = os.path.join(dst_root, name)
if os.path.exists(dst):
if os.path.getsize(dst) == os.path.getsize(src):
continue
os.remove(dst)
if name == "modular_model_index.json":
with open(src, "r", encoding="utf-8") as f:
index = json.load(f)
for value in index.values():
if isinstance(value, list) and len(value) == 3 and isinstance(value[2], dict):
if "pretrained_model_name_or_path" in value[2]:
value[2]["pretrained_model_name_or_path"] = view_dir
with open(dst, "w", encoding="utf-8") as f:
json.dump(index, f, indent=1)
continue
try:
os.link(src, dst)
except OSError:
shutil.copyfile(src, dst)
def to_audio_array(audio):
"""Normalize the pipeline output (torch tensor OR numpy array, any of
(T,), (ch, T), (T, ch), (batch, ch, T)) to float32 (channels, samples)."""
try:
import torch
if isinstance(audio, torch.Tensor):
audio = audio.detach().float().cpu().numpy()
except ImportError:
pass
import numpy as np
data = np.asarray(audio, dtype=np.float32)
if data.ndim == 3:
data = data[0]
if data.ndim == 1:
data = data[None, :]
# Channels are the small axis (stereo); samples the long one.
if data.shape[0] > data.shape[1]:
data = data.T
return data
def write_wav_i16(path, data, sample_rate):
"""data: float32 array shaped (channels, samples) in [-1, 1]."""
import wave
import numpy as np
pcm = (np.clip(data.T, -1.0, 1.0) * 32767.0).astype("<i2")
with wave.open(path, "wb") as w:
w.setnchannels(pcm.shape[1])
w.setsampwidth(2)
w.setframerate(int(sample_rate))
w.writeframes(pcm.tobytes())
def main():
parser = argparse.ArgumentParser()
parser.add_argument("--model-dir", required=True)
parser.add_argument("--view-dir", required=True)
args = parser.parse_args()
ev(ev="stage", stage="boot")
ev(ev="stage", stage="build-view")
build_local_view(args.model_dir, args.view_dir)
ev(ev="stage", stage="load-libs")
import torch # noqa: E402
from diffusers import ModularPipeline # noqa: E402
try:
ev(ev="stage", stage="load-components")
offload = os.environ.get("MAKEPAD_MUSIC3_OFFLOAD") == "1"
if offload:
# The model card's low-VRAM path: auto CPU offload keeps peak
# under ~22 GB at the cost of per-job PCIe restreaming.
from diffusers import ComponentsManager
manager = ComponentsManager()
manager.enable_auto_cpu_offload(device="cuda")
pipe = ModularPipeline.from_pretrained(
args.view_dir, components_manager=manager
)
pipe.load_components(dtype=torch.bfloat16)
else:
pipe = ModularPipeline.from_pretrained(args.view_dir)
pipe.load_components(dtype=torch.bfloat16)
ev(ev="stage", stage="to-gpu")
pipe.to("cuda")
except Exception as e: # noqa: BLE001
import traceback
traceback.print_exc()
ev(ev="error", message="load failed: %s" % str(e)[:400])
return 1
# Progress: the long pole is the 8B global LM generating audio frames at
# 25 frames/s of song. MiniMaxMusic3SemanticGenerationStep deliberately
# calls `language_model.model(...)`, not the outer Qwen3ForCausalLM
# wrapper, so the hook MUST sit on that inner model. Hooking the wrapper
# looks plausible but produces no events and leaves clients parked at the
# initial "generate" fraction for most of a song.
#
# The flow tail invokes the transformer twice per scheduler step (CFG),
# and the vocoder once per audio chunk. Their totals are derived lazily
# from the number of semantic forwards actually completed, so early EOS
# songs still report useful, bounded work progress. These are completed
# model forwards, not a wall-clock animation.
state = {
"lm": 0,
"lm_n": 0,
"dit": 0,
"dit_n": 0,
"vocoder": 0,
"vocoder_n": 0,
}
def chunk_count_from_semantic_forwards():
# One inner-model call is prompt prefill. With a full-length result the
# remaining calls equal the emitted frames; early EOS can make this
# estimate one frame high, which does not affect a chunk except exactly
# on a 100-frame boundary and is safely clamped by the Rust consumer.
frames = max(1, min(state["lm_n"], state["lm"] - 1))
return 1 if frames <= 200 else len(range(0, frames - 100, 100))
def hook_forward(module, stage_name, every, n_key=None):
orig = module.forward
def counted(*a, **kw):
result = orig(*a, **kw)
key = stage_name
state[key] += 1
if key == "dit" and state["dit_n"] == 0:
chunks = chunk_count_from_semantic_forwards()
# Official pipeline default: 30 flow steps, two CFG forwards.
state["dit_n"] = chunks * 30 * 2
state["vocoder_n"] = chunks
if state[key] % every == 0:
n = state.get(n_key, 0) if n_key else 0
ev(
ev="stage",
stage="%s %d/%d" % (key, state[key], n) if n else key,
k=state[key],
n=n if n else None,
)
return result
module.forward = counted
return orig
try:
lm = getattr(pipe, "language_model", None)
if lm is not None:
# The pinned official semantic block bypasses Qwen3ForCausalLM's
# forward and calls its `.model` core directly.
hook_forward(getattr(lm, "model", lm), "lm", 5, n_key="lm_n")
dit = getattr(pipe, "transformer", None)
if dit is not None:
hook_forward(dit, "dit", 5, n_key="dit_n")
vocoder = getattr(pipe, "vocoder", None)
if vocoder is not None:
hook_forward(vocoder, "vocoder", 1, n_key="vocoder_n")
except Exception: # noqa: BLE001
pass
ev(ev="ready")
for line in sys.stdin:
line = line.strip()
if not line:
continue
try:
job = json.loads(line)
except ValueError:
ev(ev="error", message="bad job line")
continue
if job.get("exit"):
break
state["lm"] = 0
state["dit"] = 0
state["dit_n"] = 0
state["vocoder"] = 0
state["vocoder_n"] = 0
try:
duration = float(job.get("duration_s", 60.0))
prompt = str(job.get("prompt", "")).strip()
lyrics = str(job.get("lyrics", "")).strip() or "[Instrumental]"
if not prompt:
raise ValueError("music description must not be empty")
# The released checkpoint is 25 Hz, but use the loaded component's
# authoritative frame rate so progress survives compatible model
# revisions without changing the wire protocol.
state["lm_n"] = max(1, int(duration * float(pipe.frame_rate)))
out_wav = job["out_wav"]
os.makedirs(os.path.dirname(out_wav), exist_ok=True)
ev(ev="stage", stage="generate")
generator = torch.Generator("cuda").manual_seed(int(job.get("seed", 0)))
audio = pipe(
prompt=prompt,
lyrics=lyrics,
audio_duration=duration,
generator=generator,
output="audios",
)[0]
data = to_audio_array(audio)
# Vocoder config pins 44100; pipe.sampling_rate wins when present.
sample_rate = int(getattr(pipe, "sampling_rate", 44100) or 44100)
ev(
ev="stage",
stage="write sr=%d ch=%d samples=%d"
% (sample_rate, data.shape[0], data.shape[1]),
)
write_wav_i16(out_wav, data, sample_rate)
if not os.path.isfile(out_wav):
ev(ev="error", message="no wav produced")
continue
ev(ev="done", wav=out_wav)
except Exception as e: # noqa: BLE001
import traceback
traceback.print_exc()
ev(ev="error", message=str(e)[:400])
return 0
if __name__ == "__main__":
sys.exit(main())