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>
280 lines
11 KiB
Python
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())
|