makepad/libs/ai/hub/python/fw_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

146 lines
5.3 KiB
Python

# fw_worker.py - persistent FlashWorld job worker for the makepad-ai-content
# `world` domain backend (world_backend.rs). Line protocol:
# stdin : one JSON object per line: {"prompt": str, "image_path": str|absent,
# "seed": int, "out_dir": 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","ply":path} job finished
# {"ev":"error","message":text} load/job failed (worker lives on
# after job errors)
# Runs with cwd = the FlashWorld repo clone so `cli`/`app` import; the repo's
# app.py is patched box-side to emit @EV load-stage events from
# GenerationSystem.__init__ (patch_fw2.py).
import argparse
import copy
import functools
import json
import os
import sys
def ev(**kw):
sys.stdout.write("@EV " + json.dumps(kw) + "\n")
sys.stdout.flush()
def main():
parser = argparse.ArgumentParser()
parser.add_argument("--ckpt", required=True)
parser.add_argument("--cameras", required=True)
args = parser.parse_args()
ev(ev="stage", stage="boot")
sys.path.insert(0, os.getcwd())
with open(args.cameras, "r") as f:
preset = json.load(f)
n_frame, image_height, image_width = preset["resolution"]
ev(ev="stage", stage="load-libs")
import torch # noqa: E402
from PIL import Image # noqa: E402
# cli imports app (GenerationSystem) and gsplat; module level only defines.
from cli import process_generation_request # noqa: E402
from app import GenerationSystem # noqa: E402
try:
system = GenerationSystem(ckpt_path=args.ckpt, device=torch.device("cuda"))
except Exception as e: # noqa: BLE001
import traceback
traceback.print_exc()
ev(ev="error", message="load failed: %s" % str(e)[:400])
return 1
# Denoise progress: the DiT runs exactly len(denoising_steps) forwards per
# job (3 feedback steps + 1 final). functools.wraps keeps the signature
# visible to any inspect-based callers (the H3 lesson).
n_steps = len(system.denoising_steps)
state = {"k": 0}
orig_forward = system.transformer.forward
@functools.wraps(orig_forward)
def counted_forward(*a, **kw):
state["k"] += 1
k = min(state["k"], n_steps)
ev(ev="stage", stage="denoise %d/%d" % (k, n_steps), k=k, n=n_steps)
return orig_forward(*a, **kw)
system.transformer.forward = counted_forward
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["k"] = 0
try:
out_dir = job["out_dir"]
os.makedirs(out_dir, exist_ok=True)
image_path = job.get("image_path")
if image_path:
# Normalize to the preset's exact canvas so the embedded
# camera intrinsics pass through cli.py's crop/rescale
# unchanged (scale == 1 path).
img = Image.open(image_path).convert("RGB")
w, h = img.size
target_aspect = image_width / image_height
if w / h > target_aspect:
new_w = int(round(h * target_aspect))
x0 = (w - new_w) // 2
img = img.crop((x0, 0, x0 + new_w, h))
else:
new_h = int(round(w / target_aspect))
y0 = (h - new_h) // 2
img = img.crop((0, y0, w, y0 + new_h))
img = img.resize((image_width, image_height), Image.LANCZOS)
norm_path = os.path.join(out_dir, "input_704x480.png")
img.save(norm_path)
image_path = norm_path
data = {
"text_prompt": job.get("prompt", ""),
"resolution": preset["resolution"],
"image_index": preset["image_index"],
# process_generation_request mutates camera intrinsics in
# place when an image is supplied - deep-copy per job.
"cameras": copy.deepcopy(preset["cameras"]),
}
if image_path:
data["image_prompt"] = image_path
torch.manual_seed(int(job.get("seed", 0)))
ev(ev="stage", stage="generate")
result = process_generation_request(
data, system, out_dir, video=False, spz=False, ply=True
)
if isinstance(result, dict) and result.get("error"):
ev(ev="error", message=str(result["error"])[:400])
continue
ev(ev="stage", stage="export")
ply = os.path.join(out_dir, "gaussians.ply")
if not os.path.isfile(ply):
ev(ev="error", message="no gaussians.ply produced")
continue
ev(ev="done", ply=ply)
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())