Compare commits

...
Author SHA1 Message Date
will 4e1603634d [test]: regression guard for magi-human SR layout invalidation
Pure logic test (no GPU, no model load, no upstream daVinci-MagiHuman
clone) that asserts MagiHumanSRLatentPreparationStage.forward() clears
batch.magi_static_packed_layout. Catches the f1eeb630 regression that
the existing SR-540p pipeline-parity test missed because that test
bypasses the production stage composition and re-implements the SR
denoise loop with simplified inline helpers, so the cross-stage state
transfer via ForwardBatch is never exercised.

Verified the test fails against HEAD~1 (without the f1eeb630 fix) with:
  AssertionError: assert <sentinel object> is None
and passes against HEAD.
2026-05-05 15:57:29 -07:00
will f1eeb6303f [fix]: magi-human SR latent prep: invalidate stale packed layout
C4 (4190c720) added a precompute_static_packed_layout call in the base
latent prep stage that stashes coords/modality-maps/max_ch on
batch.magi_static_packed_layout, sized for the BASE-resolution latent.

The SR latent prep stage upsamples batch.latents to the SR grid (e.g.
256x480 -> 512x896 for SR-540p), changing video_token_num and the
shapes of video_coords / video_mm — but it didn't invalidate the
precomputed layout. The SR denoising loop then passed the stale
base-sized layout to build_static_packed_inputs, which produced a
modality_mapping whose first-dim mismatched the SR-sized token tensor,
crashing in MagiHumanDiT.adapter at the text_mask scatter:

  IndexError: The shape of the mask [3243] at index 0 does not match
  the shape of the indexed tensor [11771, 3584] at index 0

Fix: clear batch.magi_static_packed_layout in MagiHumanSRLatentPrep so
the SR denoising loop falls back to the slow path of
build_static_packed_inputs (which rebuilds from the current latent
shape). Base C4 perf win is preserved (32 base steps); SR has only ~5
steps so the meshgrid recompute cost is negligible.

Repro: examples/inference/basic/basic_magi_human_sr540p.py now runs
end-to-end (34s on B200). SR-540p parity tests t2v + ti2v + DiT parity
+ distill DiT parity all pass.
2026-05-05 15:47:43 -07:00
will 990d2c2410 [feat]: magi-human DiT: enable FLASH_ATTN backend + use it for SSIM
The MagiAttention LocalAttention layer was hardcoded to TORCH_SDPA. The
attention dispatch already routes through FastVideo's selector, so
adding FLASH_ATTN to the supported list lets bf16 inference paths pick
up Hopper FA-3/FA-4 (or FA-2 elsewhere) automatically. Falls back to
SDPA for fp32 (parity tests) which the selector handles cleanly.

Switch the SSIM test to FLASH_ATTN since that's the production
inference backend; pin the parametrize list to a single backend so the
seeded reference videos correspond to the path users actually run.

All 8 runnable magi-human parity tests still pass under
FASTVIDEO_ATTENTION_BACKEND=FLASH_ATTN.
2026-05-05 14:09:11 -07:00
will 1772226bf1 [ci]: magi-human SSIM: align with example preset (width=480, steps=8)
Match basic_magi_human.py code path exactly except for CI budget knobs:
- width 448 -> 480 (preset default, was a typo'd test value)
- num_inference_steps 4 -> 8 (4 was too low to produce a stable
  regression baseline; 8 mirrors the distill preset and gives a
  recognizable but cheap sample)
- All other knobs (height, guidance_scale, seed, fps, model_path,
  cfg_number from pipeline_config, negative_prompt from preset) match
  the registered magi_human_base preset and basic_magi_human.py.
2026-05-05 12:40:18 -07:00
will 1fe8e64e23 [ci]: magi-human SSIM: bump REQUIRED_GPUS=2 + sp_size=2 for L40S fit
Single L40S (44 GB) OOMs during MagiHuman model load: 15B DiT +
T5-Gemma 9B encoder + Wan2.2 VAE + Stable-Audio VAE total ~56 GB bf16.
FSDP across 2 L40S shards the DiT + text encoder so each rank stays
under 44 GB. Mirrors test_ltx2_similarity.py and test_wan_t2v_similarity.py
which use the same 2-GPU FSDP layout for 5-15B-class models on L40S.
2026-05-05 11:49:17 -07:00
will 46f18d8cd2 [fix]: magi-human latent prep: honor batch.num_frames
The latent preparation stage was deriving num_frames from
batch.num_seconds * fps + 1 unconditionally and ignoring
batch.num_frames. num_seconds was never set anywhere (no plumbing in
SamplingParam or ForwardBatch), so every generation defaulted to 4
seconds = 101 frames, regardless of what the caller passed via
SamplingParam.num_frames.

Now we prefer batch.num_frames when it's a sensible video length (>1)
and fall back to the seconds-based derivation when the caller explicitly
passes num_seconds. Backward-compatible with the production preset
(num_frames=101 == 4*25+1). Fixes the SSIM test budget — it was running
at 101 frames despite asking for 26.

All 8 runnable magi-human parity tests still pass.
2026-05-05 11:38:17 -07:00
will 2e8db18d94 [ci]: magi-human SSIM test: switch model path to umbrella scheme
Aligns the SSIM test with the rest of the magi-human ports which now use
FastVideo/MagiHuman-Diffusers/base via maybe_download_model's umbrella-
repo support (org/repo/subfolder). The standalone -Base-Diffusers repo
no longer exists publicly; the umbrella repo is the canonical source.
2026-05-05 11:15:08 -07:00
will 4190c7203f [perf]: magi-human denoising: precompute step-invariant packed layout
build_static_packed_inputs was called every step inside the denoising
loop, redoing the meshgrid/torch.full() coords + modality-map work even
though those depend only on latent shape, audio length, and channel
widths — all fixed for a single generation.

Add StaticPackedLayout + precompute_static_packed_layout. The latent prep
stage stashes the layout on the batch; the denoise / sr-denoise loops
pass it back through the new layout= arg, which short-circuits the
invariant work and only rebuilds the per-step token tensors. The slow
path (layout=None) is kept bit-exact for build_packed_inputs callers in
parity tests.

Bit-exact verified slow vs fast path. dit / distill_dit / pipeline_smoke
/ sr540p / sr1080p / vae parity tests pass.
2026-05-05 11:11:41 -07:00
will 8f1443f47b [fix]: magi-human audio decode: scipy.signal.resample for upstream parity
Replace F.interpolate(mode='linear') with scipy.signal.resample to match
upstream video_process.resample_audio_sinc which uses the same FFT-based
polyphase resampler. Removes high-frequency aliasing and roll-off the
linear path introduced. scipy is already a direct fastvideo dep.
2026-05-05 11:11:30 -07:00
will 42ed546a66 [fix]: magi-human pipeline: don't clobber bundled components on lazy-load
Pre-detect bundled state by reading model_index.json upfront and only
defer-remove non-bundled keys from required_config_modules. After
super().load_modules, prefer modules.get(key) over loaded_modules so
super-loaded bundled or caller-provided overrides aren't silently
clobbered with a fresh upstream lazy-load.
2026-05-05 10:55:49 -07:00
SolitaryThinker d6e020402e [feat] magi-human: switch all 8 examples to umbrella HF repo + register umbrella paths
Phase 5: with all 4 weight variants now uploaded under
FastVideo/MagiHuman-Diffusers (one HF repo, four sibling subfolders
base / distill / sr_540p / sr_1080p), all magi-human examples now
default to the umbrella string. Local conversion via the
checkpoint_conversion script remains supported and is documented in
the example docstrings.

Files updated:
  * basic_magi_human.py:                      base T2V    -> base
  * basic_magi_human_ti2v.py:                 base TI2V   -> base
  * basic_magi_human_distill.py:              distill T2V -> distill
  * basic_magi_human_distill_ti2v.py:         distill TI2V-> distill
  * basic_magi_human_sr540p.py:               sr-540p T2V -> sr_540p
  * basic_magi_human_sr540p_ti2v.py:          sr-540p TI2V-> sr_540p
  * basic_magi_human_sr1080p.py:              sr-1080p T2V-> sr_1080p
  * basic_magi_human_sr1080p_ti2v.py:         sr-1080p TI2V-> sr_1080p

  * fastvideo/registry.py: hf_model_paths for the base T2V config now
    includes 'FastVideo/MagiHuman-Diffusers/base' and the distill T2V
    config includes 'FastVideo/MagiHuman-Diffusers/distill', alongside
    the existing per-variant repo names. The SR umbrella paths
    (sr_540p / sr_1080p) were already registered by Phases 3/4. TI2V
    variants reuse the same weight subfolders and are selected at
    load time via override_pipeline_cls_name + pipeline_config (see
    basic_magi_human_ti2v.py for the pattern).

Verified end-to-end:
  * pytest tests/local_tests/magi_human/ -v -s: 14 passed, 0 failed.
    All bit-exact (diff_max=0.0, diff_mean=0.0).
  * basic_magi_human.py with local converted_weights/magi_human_base
    moved out, forcing snapshot_download from the umbrella repo:
    mp4 byte-identical md5 dcf5f2bf6534c7c0d91e7353e42b23db (matches
    pre-upload local-path output exactly).

Notes:
  * fastvideo/utils.py:maybe_download_model already supports the
    'org/repo/subfolder' umbrella form (committed in e2ef3234), and
    fastvideo/pipelines/basic/magi_human/magi_human_pipeline.py
    lazy-loads the four shared components (Wan VAE, T5-Gemma encoder
    + tokenizer, Stable Audio VAE) from their canonical upstream HF
    repos so each umbrella subfolder only ships transformer/ +
    scheduler/ + (sr_transformer/) + model_index.json.
  * Total HF Hub footprint after upload: ~165 GB raw, server-side
    deduped. User cache footprint per variant is ~5-30 GB transformer
    (+ ~30 GB sr_transformer for SR) plus a single ~25 GB upstream
    cache shared across all variants.
2026-05-05 10:55:49 -07:00
SolitaryThinker 620e100af4 [feat] magi-human: SR-1080p with block-sparse local-window attention
Ports the daVinci-MagiHuman SR-1080p inference flow to FastVideo. Builds
on the SR-540p two-stage pipeline plus block-sparse local-window
video->video attention on 32 of 40 SR DiT layers. Mirrors upstream
SR2_1080 config override at inference/common/config.py:225-244 and
calc_local_qk_range at inference/pipeline/data_proxy.py:31-79.

Files added:
  * tests/local_tests/magi_human/test_magi_human_sr1080p_pipeline_parity.py:
    parametric parity test for both T2V and TI2V modes. Both pass
    diff_max=0.0000 / diff_mean=0.0000 -- BIT-EXACT.

Files modified (key changes):
  * fastvideo/models/dits/magi_human.py:
    - AttentionSubConfig.use_local_attn / frame_receptive_field
    - MagiAttention.configure_local_attention(): per-layer toggle
    - MagiAttention._sdpa(): thin SDPA wrapper for [L,H,D] tensors
    - MagiAttention._local_window_attention(): 3-block accumulator
      mirroring upstream FFA semantics with vanilla SDPA segments
      (per-frame video window + all video->audio+text + audio+text->all).
    - MagiAttention.forward() dispatches to _local_window_attention
      when use_local_attn flag is set.
    - MagiTransformerLayer wires use_local_attn from arch.local_attn_layers.
    - MagiHumanDiT.configure_local_attention() top-level toggle.
  * pipeline_configs.py: MagiHumanSR1080pConfig + I2V variant with
    sr_local_attn_layers populated to upstream's 32 indices.
  * presets.py: MAGI_HUMAN_SR_1080P + MAGI_HUMAN_SR_1080P_TI2V presets.
  * registry.py: SR-1080p config entries with detectors.
  * magi_human_pipeline.py: SR-1080p pipeline classes activating
    local_attn_layers on the SR DiT at construction.
  * scripts/checkpoint_conversion/convert_magi_human_to_diffusers.py:
    --sr-subfolder 1080p_sr support.
  * tests/local_tests/helpers/magi_human_upstream.py: arch override
    support for SR-1080p parity test.
  * test_magi_human_pipeline_smoke.py: preset set expanded.
  * basic_magi_human_sr1080p{,_ti2v}.py: stubs -> runnable.

Verification:
  * pytest tests/local_tests/magi_human/ -v -s: 14 passed, 0 failed.
    All bit-exact including SR-1080p T2V and TI2V parity. The
    block-sparse SDPA-segmented implementation matches upstream FFA's
    q_ranges/k_ranges accumulator semantics exactly for the 3-block
    layout (overlap accumulation handled by explicit '+' at
    magi_human.py:421-425).
  * basic_magi_human.py (T2V regression check): mp4 byte-identical
    md5 dcf5f2bf6534c7c0d91e7353e42b23db.
  * Base/Distill/TI2V/SR-540p flows untouched.

Notes:
  * Converted SR-1080p artifact at /raid/.../magi_human_sr_1080p
    (~58 GB), symlinked at converted_weights/magi_human_sr_1080p
    (gitignored).
  * Upstream's flex_flash_attn_func via SandAI-org/MagiAttention is
    NOT a dependency; FV's pure-SDPA segmented implementation is
    mathematically equivalent for this 3-block layout.
2026-05-05 10:55:49 -07:00
SolitaryThinker 0fde316a19 [feat] magi-human: SR-540p two-stage super-resolution pipeline
Ports the daVinci-MagiHuman SR-540p inference flow to FastVideo for
both T2V and TI2V modes. SR-540p is a TWO-STAGE pipeline: the base
model produces a 256x480 latent, then a separate SR DiT (same arch as
base, different weights) refines it to 896x512. Mirrors upstream
MagiEvaluator.evaluate at video_generate.py:300-360.

Files added:
  * fastvideo/pipelines/basic/magi_human/stages/sr_latent_preparation.py
  * fastvideo/pipelines/basic/magi_human/stages/sr_denoising.py
  * tests/local_tests/magi_human/test_magi_human_sr540p_pipeline_parity.py

Files modified (key changes):
  * pipeline_configs.py: SR config classes with sr_* knobs sourced
    from upstream EvaluationConfig
  * presets.py: MAGI_HUMAN_SR_540P + MAGI_HUMAN_SR_540P_TI2V
  * registry.py: SR-540p config entries
  * magi_human_pipeline.py: MagiHumanSRPipeline + MagiHumanSRI2VPipeline
    classes wiring the 9-stage chain (base denoise -> sr latent prep
    -> sr denoise -> decode)
  * fastvideo/models/loader/component_loader.py: registers
    'sr_transformer' alongside transformer / transformer_2
  * scripts/checkpoint_conversion/convert_magi_human_to_diffusers.py:
    --sr-source / --sr-subfolder flags for SR DiT into sr_transformer/
  * examples/inference/basic/basic_magi_human_sr540p{,_ti2v}.py:
    runnable

Verification:
  * pytest tests/local_tests/magi_human/ -v -s: 12 passed, 0 failed.
    All bit-exact (diff_max=0.0, diff_mean=0.0) including new SR-540p
    T2V + TI2V parity tests.
  * basic_magi_human.py (T2V regression check): mp4 byte-identical
    md5 dcf5f2bf6534c7c0d91e7353e42b23db.
  * basic_magi_human_sr540p.py: 896x512 mp4, coherent reading-on-
    park-bench scene, much higher quality than base 480x256.
  * basic_magi_human_sr540p_ti2v.py: 896x512 mp4 with reference-image-
    conditioned saxophonist; reference conditioning preserved through
    SR upscale.

Notes:
  * Converted SR-540p artifact at converted_weights/magi_human_sr_540p
    (~86 GB; both transformer/ and sr_transformer/) is gitignored.
  * Base T2V/TI2V/distill flows untouched and bit-exact.
2026-05-05 10:55:49 -07:00
SolitaryThinker 0d99e47e16 [feat] magi-human: TI2V (text+image-to-AV) inference flow
Ports the daVinci-MagiHuman TI2V branch to FastVideo for both base
and distill variants. The TI2V case takes a reference image, encodes
it through the Wan VAE, and overwrites the first frame's video latent
with the encoded image latent at every denoise step (mirrors upstream
inference/pipeline/video_generate.py:300-360 evaluate + 424-425
per-step overwrite).

Files added:
  * fastvideo/pipelines/basic/magi_human/stages/reference_image.py:
    new MagiHumanReferenceImageStage. Loads PIL image (or path),
    resizecrops to (height, width) matching upstream resizecrop,
    runs VideoProcessor.preprocess at vae_scale_factor=16, encodes
    via the Wan VAE (uses .mean for deterministic latent), applies
    shift_factor / scaling_factor normalization, stashes on
    batch.image_latent.
  * tests/local_tests/magi_human/test_magi_human_ti2v_pipeline_parity.py:
    bit-exact parity test against upstream MagiEvaluator's TI2V denoise
    (with reference image conditioning). Passes
    ti2v video diff_max=0.0000 diff_mean=0.0000.

Files modified:
  * pipeline_configs.py: MagiHumanBaseI2VConfig keeps the VAE encoder
    loaded (load_encoder=True) so the reference image path can encode.
  * presets.py: MAGI_HUMAN_BASE_TI2V and MAGI_HUMAN_DISTILL_TI2V presets
    with workload_type=i2v.
  * registry.py: TI2V config entries for both base and distill variants.
  * magi_human_pipeline.py: MagiHumanI2VPipeline subclass that inserts
    the MagiHumanReferenceImageStage between prompt encoding and latent
    preparation. Reuses the lazy-load path for shared components.
  * stages/latent_preparation.py: pre-loop overwrite of
    latent_video[:, :, :1] with batch.image_latent[:, :, :1] when
    image conditioning is present (matches upstream
    evaluate_with_latent line 425 first-iteration overwrite).
  * stages/denoising.py: per-step _overwrite_first_frame helper that
    applies the same overwrite at the start of every denoise step
    (matches upstream evaluate_with_latent line 424 in-loop overwrite).
    static_packed rebuild moved inside the loop after the overwrite
    so packed video tokens reflect the conditioned latent.
  * examples/inference/basic/basic_magi_human_ti2v.py: rewritten from
    NotImplementedError stub to runnable. Uses local
    converted_weights/magi_human_base + the existing example
    saxophonist reference image; produces a coherent mp4 with the
    image-conditioned subject.
  * examples/inference/basic/basic_magi_human_distill_ti2v.py: same
    pattern against converted_weights/magi_human_distill (runnable
    once distill weights are converted).

Verification:
  * pytest tests/local_tests/magi_human/ -v -s: 10 passed, 0 failed.
    All bit-exact (diff_max=0.0, diff_mean=0.0) including new
    test_magi_human_distill_dit_parity (Phase 1) and
    test_magi_human_ti2v_pipeline_parity (Phase 2) tests.
  * basic_magi_human.py (T2V regression check): mp4 byte-identical
    md5 dcf5f2bf6534c7c0d91e7353e42b23db.
  * basic_magi_human_ti2v.py: produces coherent saxophonist scene
    matching the reference image conditioning. Frames at
    /tmp/opencode/ti2v_frame_*.png.

T2V flow unchanged: both T2V example and base parity test produce
identical output to pre-change. TI2V is purely additive on top.
2026-05-05 10:55:49 -07:00
SolitaryThinker 27f6f0aacd [bugfix] convert_magi_human: keep fp32 for adapter+final_linear weights
The conversion script's _FP32_KEEP_SUFFIXES list was missing 8 keys
that the BASE checkpoint stores as fp32:

  adapter.video_embedder.{weight,bias}
  adapter.text_embedder.{weight,bias}
  adapter.audio_embedder.{weight,bias}
  final_linear_video.weight
  final_linear_audio.weight

For the BASE conversion this never surfaced because the BASE checkpoint
already ships these as fp32 (no --cast-bf16 needed; conversion was
identity for these weights). For the DISTILL conversion (which ships
ALL 331 weights as fp32 and relies on --cast-bf16 to produce a 30 GB
bf16 artifact), the omission caused these 8 fp32 layers to be cast
to bf16, which mismatched FV's MagiAdapter/final_linear modules
(declared dtype=torch.float32 at magi_human.py:519-527 and 645-648,
mirroring upstream Adapter at dit_module.py:721-723 and DiTModel at
dit_module.py:896-900).

Symptom: distill DiT parity vs upstream had video diff_mean=0.114
(19% relative error) instead of the expected 0.0. Adding the 8 keys
to _FP32_KEEP_SUFFIXES restores bit-exact parity.

Also adds tests/local_tests/magi_human/test_magi_human_distill_parity.py
mirroring the base DiT parity test but pointing at the distill shards
and converted weights. Bit-exact diff_max=0.0, diff_mean=0.0 vs
upstream daVinci-MagiHuman/inference/model/dit/dit_module.py:DiTModel
loaded from the distill subfolder.

Verified with --cast-bf16 reconversion of GAIR/daVinci-MagiHuman/distill
into converted_weights/magi_human_distill (29 GB).
2026-05-05 10:55:49 -07:00
SolitaryThinker c97fb6b3b3 [feat] utils: support umbrella-repo subfolders in maybe_download_model
Recognises an 'umbrella' HF repo layout where a single repo holds
multiple pipeline variants under sibling subfolders, e.g.

  FastVideo/MagiHuman-Diffusers/
    base/{model_index.json, transformer/, scheduler/}
    distill/{...}
    sr_540p/{...}
    sr_1080p/{...}

Users can pass 'org/repo/subfolder' as the model path and the loader
downloads only that subfolder's blobs (allow_patterns=['<sub>/**'])
and returns the local subfolder snapshot path:

  generator = VideoGenerator.from_pretrained(
      'FastVideo/MagiHuman-Diffusers/base',
  )

Detection is structural: HF Hub repo ids are always two
slash-separated components; a path with 3+ components that does not
exist locally and is not posix-absolute or relative-prefixed is
treated as an umbrella reference. The existing single-repo-per-variant
layout ('FastVideo/MagiHuman-Base-Diffusers') still works unchanged.

Combined with the lazy-load of the four cross-variant shared
components landed in 53ac1985, an umbrella MagiHuman repo only needs
to ship transformer/+scheduler/+model_index.json per variant and the
user's local cache stays at ~75 GB total for all 4 variants instead
of ~400 GB.

Documented in tests/local_tests/magi-human.md under 'Design notes'.
Verified existing converted_weights/magi_human_base local-path flow
still produces a byte-identical mp4 (md5 dcf5f2bf...) to the pre-
refactor reference.
2026-05-05 10:55:49 -07:00
SolitaryThinker 3464cb8b03 [refactor] magi-human: lazy-load Wan VAE from upstream (drop bundling default)
Extends the existing lazy-load pattern for text_encoder, tokenizer,
and audio_vae to the video VAE: each MagiHuman variant's converted
repo no longer needs to bundle a copy of the Wan 2.2 TI2V-5B VAE.
The four cross-variant shared components are now all fetched from
their canonical upstream HF repos at first build:

  * text_encoder, tokenizer  -> google/t5gemma-9b-9b-ul2 (gated)
  * audio_vae                -> stabilityai/stable-audio-open-1.0 (gated)
  * vae                      -> Wan-AI/Wan2.2-TI2V-5B-Diffusers

Per-variant converted repo shrinks to transformer/ + scheduler/ +
model_index.json (~5 GB for base bf16, ~30 GB for distill bf16). All
variants share the same ~25 GB cache of upstream weights, so a user
running 4 variants ends up with ~75 GB total instead of ~400 GB.

Implementation:

  * fastvideo/utils.py:verify_model_config_and_directory now treats
    the contents of model_index.json as authoritative for which
    component subfolders must exist locally. Pipelines that emit a
    minimal model_index.json (omitting vae / text_encoder / etc.)
    pass verification; pipelines that DO declare a component must
    still ship its subfolder. transformer/ remains mandatory.

  * fastvideo/pipelines/basic/magi_human/magi_human_pipeline.py adds
    vae to the deferred list in load_modules and a new
    _load_video_vae helper that prefers a bundled vae/ subfolder
    (legacy converted repos) and falls back to snapshot_download +
    the standard FV VAELoader. Both paths produce the same FV
    AutoencoderKLWan, so production behavior is unchanged.

  * convert_magi_human_to_diffusers.py docstring updated; --bundle-vae
    flag is unchanged (still optional) but the README example now
    omits it so new converted repos default to the minimal layout.

Verification:
  * basic_magi_human.py with bundled vae/: byte-identical mp4 (md5
    dcf5f2bf...) to pre-refactor output.
  * basic_magi_human.py with vae/ moved out and removed from
    model_index.json: same byte-identical mp4 via the lazy-load path.
  * 4/4 magi-human parity tests pass with diff_max=0.0 / diff_mean=0.0
    (DiT, pipeline, smoke + typed surface preflight).
2026-05-05 10:55:49 -07:00
SolitaryThinker e114fba53f [docs] examples: add basic_* stubs for remaining magi-human variants
Upstream daVinci-MagiHuman ships 4 model variants x 2 input modes = 8
inference entrypoints (base / distill / sr_540p / sr_1080p, each in
T2V and TI2V mode). FastVideo currently has working code for the base
T2V path (basic_magi_human.py) and a registered preset for distill
T2V (magi_human_distill, but no example until now).

Add example files for the 7 remaining variants:

  * basic_magi_human_distill.py -- runnable T2V example for the
    DMD-2 distilled model. Just point conversion at the distill
    subfolder and the existing magi_human_distill preset takes over.

  * basic_magi_human_ti2v.py
    basic_magi_human_distill_ti2v.py -- not-yet-ported TI2V (image
    conditioning) variants. Each docstring lists the pipeline-side
    work that is missing in FastVideo (VAE encoder load, reference
    image stage, latent_video[..., :1] overwrite at every denoise
    step, new I2V config + preset). main() raises NotImplementedError
    with a pointer to magi-human.md.

  * basic_magi_human_sr540p.py
    basic_magi_human_sr540p_ti2v.py -- not-yet-ported super-resolution
    to 540p. Docstring describes the upstream two-stage flow (base ->
    SR latent prep with trilinear up + ZeroSNR noise -> SR DiT -> Wan
    VAE) and lists the FV components needed (MagiHumanSR540pConfig,
    SR latent prep stage, SR denoise stage with cfg-trick guidance
    tensor, conversion script invocation for the 540p_sr subfolder,
    new preset + registry entry).

  * basic_magi_human_sr1080p.py
    basic_magi_human_sr1080p_ti2v.py -- not-yet-ported SR-1080p.
    Same SR-540p scaffolding plus a block-sparse local-window
    attention path (32 of 40 SR DiT layers). The docstring points at
    upstream's FFAHandler q_ranges/k_ranges blocks in dit_module.py
    and notes that MagiAttention currently always runs full SDPA.

All stubs follow the existing example file convention (SPDX header,
focused docstring, single main()). The not-yet-ported stubs exit with
NotImplementedError so they fail loudly rather than silently misbehaving;
each error message points at the docstring for the missing-component
checklist.
2026-05-05 10:55:49 -07:00
SolitaryThinker 093f5e699c [refactor] magi-human DiT: reuse fastvideo.layers RoPE primitive
Drop the file-local apply_rotary_emb / _rotate_half helpers and call
fastvideo.layers.rotary_embedding._apply_rotary_emb with
is_neox_style=True instead. Magi uses partial RoPE (rotate first
6 * (head_dim // 8) = 96 of 128 head_dim positions, leave the trailing
32 unrotated), which the FV primitive does not handle directly, so
the partial-RoPE slicing stays in the call site:

  q_rot = _apply_rotary_emb(q[..., :rot_dim], cos, sin, is_neox_style=True)
  q = torch.cat([q_rot, q[..., rot_dim:]], dim=-1)

The math is identical to the previous local impl (Magi's
'rotate_half + doubled cos/sin' expands to upstream's 'chunk + cat(o1,
o2)' Neox formulation). Drops the einops dependency from this file.

Bit-exact preservation verified at both scales:
  * All 7 parity tests still pass with 0.0/0.0 diff vs upstream.
  * Production E2E mp4 is byte-identical (same md5 hash) to the
    post-LocalAttention output, so basic_magi_human.py output is
    unchanged.
2026-05-05 10:55:49 -07:00
SolitaryThinker e24bc12c59 [refactor] magi-human: route DiT attention through FV's LocalAttention
Replace bare F.scaled_dot_product_attention call inside MagiAttention
with FastVideo's LocalAttention layer so the backend selection (SDPA /
FlashAttn / SLA / SageAttn) flows through the standard configurable
path. Also drops the manual GQA repeat_interleave: SDPABackend's
enable_gqa=True handles num_heads_q != num_heads_kv directly.

LocalAttention requires a forward_context, so add set_forward_context
wrapping at the two call sites:

  * Production: MagiHumanDenoisingStage's per-step DiT calls now run
    inside set_forward_context(current_timestep=t, attn_metadata=None),
    matching the pattern used by the generic DenoisingStage and other
    custom denoise stages (ltx2, stable_audio, longcat).

  * Parity tests: test_magi_human_parity.py and
    test_magi_human_pipeline_parity.py wrap their direct DiT calls in
    the same context so LocalAttention's get_forward_context() doesn't
    assert.

Parity tests stay bit-exact (all 7 magi-human parity tests pass with
0.0/0.0 diff vs upstream). Production E2E still produces a coherent
video matching the prompt; the mp4 output is no longer byte-identical
to the pre-refactor version because SDPA dispatches to a different
kernel when enable_gqa=True at production sequence length (~3000
tokens), with mean per-pixel drift ~3/255 (~1.3%) -- visually
indistinguishable, just a different bf16 quantization noise pattern.

The MagiAttention block stays as a custom nn.Module rather than being
fully replaced by LocalAttention because of the per-modality packed
linears (PackedExpertLinear), per-head sigmoid gating, and per-modality
RMSNorm pattern that the FV layer abstractions don't model. The bare
SDPA call inside it is now the only piece reused from FV.
2026-05-05 10:55:49 -07:00
SolitaryThinker 4d04c1b01c [refactor] magi-human DiT: restore dtype-agnostic attention via orig_dtype
Wave 14b's bit-exact dtype boundary fix introduced two violations of
Wave 10's dtype-agnostic pattern in MagiAttention:

  1. q/k/v.to(torch.bfloat16) hardcoded before SDPA
  2. out.float() explicit upcast after attention to fp32 the gating

Both turn out to be unnecessary:

  1. orig_dtype = self.linear_qkv.weight.dtype already evaluates to
     bf16 in production (loader sets bf16) AND in the parity test
     (PackedExpertLinear's __init__ default is bf16, matching upstream
     BaseLinear at dit_module.py:330). So q.to(orig_dtype) gives the
     same bf16 cast as the hardcode without the dtype lock-in.

  2. PyTorch's float type-promotion rules already handle the
     bf16 -> fp32 transition implicitly: bf16_out * sigmoid(fp32_g)
     promotes to fp32, exactly matching upstream's intentional
     bf16 * fp32 -> fp32 boundary at dit_module.py:649. The explicit
     .float() was redundant.

Result: 11 insertions, 14 deletions, attention block reads more
cleanly without sacrificing any of the parity work:

  - All 7 magi-human local parity tests pass with bit-exact (0.0/0.0)
    or near-bit-exact (VAE 8e-4) outputs vs upstream
  - basic_magi_human.py E2E produces byte-identical mp4 (same md5
    hash) as the pre-refactor Wave 14b state, so production
    behavior is unchanged

The MagiAttention block stays custom (rather than reusing FV's
LocalAttention) because LocalAttention requires a forward_context
that the parity test does not establish; switching would force
either expanded test scaffolding or a parity-validation regression.
2026-05-05 10:55:49 -07:00
SolitaryThinker 3fb150bbe0 [bugfix] magi-human: align DiT dtype boundaries with upstream for bit-exact parity
Three cumulative dtype-boundary divergences caused FV parity tests
to fail after the Wave 14 channel-major fix exposed real-signal
processing (vs the prior garbage-in-garbage-out kernel cancellation):

1. MagiAttention forward: hardcode bf16 for SDPA q/k/v inputs to
   match upstream flash_attn_with_cp's hardcoded bf16 cast at
   dit_module.py:508. Upcast attention output to fp32 before the
   per-head gating multiply to match upstream's bf16*fp32 promotion
   at dit_module.py:649. Cast back to bf16 only for linear_proj.
   This is an intentional exception to Wave 10's dtype-agnostic
   pattern: upstream is not dtype-agnostic at this boundary.

2. MagiHumanDiT forward: drop the x.to(linear_qkv.weight.dtype)
   cast that ran the residual stream in bf16 across all 40 layers.
   Upstream casts to params_dtype which defaults to fp32, keeping
   the cross-layer accumulator fp32 with bf16 internal compute.
   The bf16 residual cast was compounding ~6-7 bits of mantissa
   loss per layer and was visible in pipeline parity diff.

3. Parity test _build_fastvideo_schedulers: switch to single-shift
   construction (default shift=1 in __init__, then set_timesteps
   with the real shift) to match production migration in Wave 11.
   The stale double-shift was applying a non-trivial shift twice
   versus upstream's single-shift, leaving FV with a different
   timestep schedule.

Result: all 7 magi-human local parity tests pass with bit-exact
or near-bit-exact (VAE 8e-4 max) outputs vs upstream:

  test_magi_human_dit_parity                     diff=0.0
  test_magi_human_t5gemma_parity                 diff=0.0
  test_magi_human_sa_audio_parity                diff=0.0
  test_magi_human_sa_audio_official_parity       diff=0.0
  test_magi_human_vae_parity                     max=8e-4
  test_magi_human_pipeline_smoke                 passes
  test_magi_human_pipeline_latent_parity         diff=0.0

Production E2E re-validated: examples/inference/basic/basic_magi_human.py
still produces coherent video at the standard 480x256 / 32-step /
seed-42 prompt, with unchanged 23.5s runtime.

Closes OQ-6 (full resolution including dtype boundaries).
2026-05-05 10:55:49 -07:00
SolitaryThinker 320f8a1f8d [bugfix] magi-human: pack video tokens channel-major to match upstream
FV's _img2tokens was packing video latent patches as spatial-major
'(pT pH pW C)' (channels innermost), but the DiT's video_embedder
Linear weight was trained on upstream's channel-major '(C pT pH pW)'
layout. Upstream's MagiDataProxy uses UnfoldNd, which is implemented
via a grouped-channel conv whose output reshape is channel-major
(in_channels * kernel_size_numel ordering, channels slowest). The
spatial-major layout silently permuted the in-features of every
video token, scrambling the entire video feature representation
and producing pure colorful-blob noise from the basic example.

The pipeline parity test could not catch this because it imports
FV's build_packed_inputs for the upstream side too, so both sides
consumed equally-permuted tokens and agreed on garbage at production
scale (T=26 frames, ~3120 video tokens), while reporting healthy
~0.5%/step drift on the tiny synthetic (2,6,6) latent. Fix is a
one-character change in the einops rearrange pattern.

unpack_tokens stays spatial-major because the DiT's final_linear_video
is trained to emit (pT pH pW C), matching upstream's
SingleData.depack_token_sequence at data_proxy.py:220-228.

Validated end-to-end: examples/inference/basic/basic_magi_human.py
now produces coherent video at the standard 480x256 / 32-step / seed
42 prompt -- woman on a park bench reading a book under green trees,
matching prompt. Output mp4 size dropped from ~932KB (incompressible
noise) to ~222KB (coherent compressible video).

Closes magi-human OQ-6.

Notes for follow-up:
- OQ-11 (NEW): test_magi_human_pipeline_parity.py:222 should drive
  the upstream side through real MagiDataProxy.process_input so the
  parity test catches packing-layout regressions natively.
2026-05-05 10:55:49 -07:00
SolitaryThinker ccfcc3042b [docs] magi-human: Wave 10 dtype refactor + post-rebase test paths + OQ-9
- Update all test path references from old bucket dirs
  (transformers/, encoders/, vaes/, pipelines/) to the post-rebase
  consolidated tests/local_tests/magi_human/ directory.
- Append Wave 10 subsection to Numerical-alignment investigation:
  WanVideo-pattern dtype refactor, 7 hardcoded bf16 casts removed,
  fp32 parity improvement, upstream residual drift explanation.
- Add OQ-9 to open questions: upstream dit_module.py hardcoded bf16
  casts block full fp32 parity validation (low-priority follow-up).
- Update Phase 11 status: Last verified 2026-05-01 Wave 10 @ 3caeaad1,
  note tests now under tests/local_tests/magi_human/.
2026-05-05 10:55:49 -07:00
SolitaryThinker 8f0493637e [refactor] magi-human DiT: remove hardcoded bf16 casts; loader-owned dtype
Match the canonical FastVideo dtype pattern (per WanVideo DiT). Remove
all 7 hardcoded `.to(torch.bfloat16)` casts that were verbatim copies of
upstream daVinci-MagiHuman/inference/model/dit/dit_module.py. Replace
with `orig_dtype = self.linear_qkv.weight.dtype` (or equivalent
loader-owned dtype) + `.to(orig_dtype)` pattern, mirroring
fastvideo/models/dits/wanvideo.py:325-392.

Refactored sites: attention pre_norm output, q/k/v post-RoPE, attention
output, MLP pre_norm, MLP activation output, top-level block-input cast.

Model dtype is now loader-owned via pipeline_config.dit_precision, not
hardcoded. Production bf16 behavior is bit-identical (loader sets
default_dtype=bf16, all params/inputs naturally bf16, orig_dtype=bf16,
outputs preserved as bf16). fp32 parity now works end-to-end on the FV
side; remaining drift in fp32 parity tests is upstream's still-hardcoded
bf16 casts (tracked as OQ-9 in tests/local_tests/magi-human.md).
2026-05-05 10:55:49 -07:00
SolitaryThinker f5ce12f17a [docs] activation trace mode + Extensions 1-3 design + magi-human Wave 9 entry
Add docs/contributing/activation_trace.md: a 210-line contributor guide
covering the env-gated activation trace infrastructure (Extension 0),
its five configuration env vars, the trace_step() context manager, and
the JSONL output format.

The doc also sketches Extensions 1-3 (FX graph capture, AST-level
instrumentation, dispatch-table interception) as future work, giving
contributors a clear design ladder to climb without requiring them to
implement everything at once.

Update tests/local_tests/magi-human.md with the Wave 9 investigation
entry: records the E2E smoke result (28,864 JSONL records, 451 hooked
modules), the zero-overhead-when-off confirmation, and the new
fastvideo/tests/hooks/test_activation_trace.py test row in the Phase 11
status table. Last-verified date bumped to 2026-05-01 (Wave 9).
2026-05-05 10:55:49 -07:00
SolitaryThinker 17cb6737c1 [feat] magi-human denoising: wire trace_step around DiT calls
Wire the new trace_step(step_idx) context manager around each DiT
forward call in the magi-human denoising stage so that per-step
activation dumps are correctly indexed.

The context manager sets a thread-local step index that hook callbacks
read when deciding whether to emit a record (controlled by
FASTVIDEO_TRACE_STEPS). Without this wiring, all records would carry
step_idx=None and step-filtered traces would be empty.

The change is a no-op when FASTVIDEO_TRACE_ACTIVATIONS is unset: the
context manager is a lightweight nullcontext in that path.
2026-05-05 10:55:49 -07:00
SolitaryThinker f39dbe482c [feat] add zero-overhead-when-off activation trace mode (env-gated)
Add an env-gated activation trace mode that registers PyTorch forward
hooks on selected modules, computes per-tensor stats (abs_mean, sum,
min, max, mean, std, shape, dtype), and writes JSONL records to a
configurable sink path.

Designed for parity-debug across model ports: enable trace on both
FastVideo's path and the upstream reference, then `diff` the two JSONL
files to find the first divergent layer. Inspired by SGLang's
--debug-tensor-dump-output-folder pattern and TransformerEngine's
DumpTensors selective-dump infra.

Zero-overhead-when-off guarantee: the master toggle
FASTVIDEO_TRACE_ACTIVATIONS is checked ONCE at pipeline startup. When
unset/false, attach_activation_trace() returns None and no hooks are
ever registered. The production forward path is untouched.

Configuration via 5 env vars (FASTVIDEO_TRACE_LAYERS regex filter,
FASTVIDEO_TRACE_STATS, FASTVIDEO_TRACE_OUTPUT path with <pid> templating,
FASTVIDEO_TRACE_STEPS step-index filter). Step indexing via the
trace_step(step_idx) context manager wires step into thread-local for
hook callbacks.

Six unit tests cover off/on/filter/stats/step-filter/tuple-flattening
/cleanup paths. End-to-end smoke against magi-human's basic example
generated 28,864 JSONL records on 451 hooked modules; trace OFF does
not create the output file.
2026-05-05 10:55:49 -07:00
SolitaryThinker 5b5608cb37 [docs] magi-human: Wave 8 broader CFG/preset/fallback audit; OQ-6 production resolved
Documents Wave 8 audit findings and targeted fixes in
tests/local_tests/magi-human.md:

- New subsection 'Wave 8 (2026-05-01)' under Numerical-alignment
  investigation: 4-item audit table (items #1, #3, #10, #7-stale),
  per-test parity numbers post-Wave-8, and OQ-6 status update.
- Phase 11 status table: pipeline parity row updated to reflect Wave 7.5
  real-prompt encoding and Wave 8 production-fix context.
- Last verified line updated to 2026-05-01 (Wave 8 broader audit + fixes).
- OQ-6 status updated to RESOLVED-PRODUCTION; bf16 amplification floor
  TRACKED separately.

Key conclusion: all production-facing root causes for OQ-6 are now
identified and fixed (incomplete neg prompt Wave 7, tokenizer pre-padding
Wave 8 #1, resolution defaults Wave 8 #3, silent audio fallback Wave 8
#10). Residual parity-test drift is the inherent bf16+CFG amplification
floor through the multistep FlowUniPC scheduler — not a code bug.
2026-05-05 10:55:49 -07:00
SolitaryThinker 1d4a6037eb [test] magi-human: encode real preset prompts in pipeline parity
test_magi_human_pipeline_parity.py previously used random txt_feat and
neg_txt_feat tensors (identical on both FV and upstream sides). This
meant the parity test never exercised the actual text-encoding path,
and production-facing preset values (the real positive/negative prompt
strings) were never validated end-to-end.

Lines 59-153 now encode the real preset prompts via T5-Gemma on both
sides before the denoising loop. Lines 388-395 wire the encoded
embeddings into the FV pipeline call. This validates that:
  - The preset prompt strings flow through T5-Gemma correctly.
  - The encoded embeddings are numerically consistent between FV and
    upstream for the same input text.
  - Production-facing preset values are exercised in the parity path.

Note: Wave 8 production fixes (tokenizer pre-padding, resolution
defaults) do NOT change parity numbers because both sides use the same
encoder/decoder. The residual drift is the inherent bf16+CFG
amplification floor (Wave 8 audit conclusion).
2026-05-05 10:55:49 -07:00
SolitaryThinker 8ac6526cdc [fix] magi-human audio decode: raise on missing audio_latents
MagiHumanAudioDecodingStage.forward() previously returned silently
(no audio output) when batch.audio_latents was None or missing. In a
joint audio-video pipeline this is a real bug: the caller expects audio
output and gets nothing, with no indication of why.

audio_decoding.py:90-96 now raises ValueError with a descriptive
message when audio_latents is absent. This converts a silent wrong
result into a loud, actionable error.

This is Wave 8 fix #10. Classified AMBIGUOUS→HARMFUL: silent return
is acceptable for optional audio in a T2V-only pipeline, but MagiHuman
is a joint AV model where missing audio_latents indicates a real
upstream failure (e.g. denoising stage dropped the audio modality).
2026-05-05 10:55:49 -07:00
SolitaryThinker 08364b2080 [fix] magi-human: align default resolution to upstream 480x256
FV's magi_human_base and magi_human_distill presets used width=448,
height=256 as defaults. Upstream daVinci-MagiHuman uses width=480,
height=272 (snapped to 256 by the latent-preparation stage). Production
users running FV with default settings got a different aspect ratio
(448/256=1.75) than upstream (480/256=1.875), causing visual composition
differences.

Fix: presets.py:50-84 now uses width=480 for both base and distill
presets. latent_preparation.py:130-133 fallback also updated to 480.

This is Wave 8 fix #3. Upstream reference: daVinci-MagiHuman/
video_generate.py default resolution args. Fixes OQ-6 production root
cause: aspect ratio mismatch between FV default and upstream default.
2026-05-05 10:55:49 -07:00
SolitaryThinker f81de3926f [fix] magi-human encoder: stop tokenizer pre-padding T5-Gemma
The T5GemmaEncoderModel._encode() call was passing truncation=True,
padding='max_length', max_length=640 to the HF tokenizer, causing the
tokenizer to pad every sequence to 640 tokens BEFORE encoding. This
means pad-token hidden states were fed into the DiT as real content,
and magi_original_text_lens reported the pre-padded length (640) rather
than the actual token count.

Upstream (daVinci-MagiHuman/models/text_encoder.py) does NOT pre-pad:
it tokenizes without padding/truncation and lets the downstream
MagiHumanLatentPreparationStage._pad_or_trim_dim1 handle length
alignment. FV now matches this: t5gemma.py:57-64 no longer passes
truncation/padding/max_length to the tokenizer call.

This is Wave 8 fix #1. Fixes OQ-6 production root cause: pad-token
hidden states were polluting DiT cross-attention input for every
inference call.
2026-05-05 10:55:49 -07:00
SolitaryThinker 04633096e2 [docs] magi-human: Wave 7 CFG/neg-prompt findings; OQ-6 partial fix
Document Wave 7 CFG + negative-prompt investigation findings in the
numerical-alignment investigation section. Key results: CFG math is
identical on both sides; audio VAE is bit-exact vs official (new parity
test confirms); root cause of production audio drift was the incomplete
negative prompt in FV's preset (missing audio-quality + speech-delivery
blocks from upstream video_generate.py:222-224).

Update Phase 11 status table with the new SA-official parity test row
(PASS, diff_max=0, diff_mean=0). Add test_magi_human_sa_audio_official_parity.py
to the running-tests command block and the What-each-test-covers section.

Mark OQ-6 as PARTIALLY-RESOLVED: production neg-prompt fix shipped in
this wave; parity-test 4-step compounding is a separate inherent
FlowUniPC scheduler phenomenon, not a code bug. Update Last-verified
to 2026-05-01 (Wave 7).
2026-05-05 10:55:49 -07:00
SolitaryThinker b6be3d0c8a [test] magi-human: add SA usage parity vs official repo (bit-exact)
Add parity test comparing FastVideo's SAAudioVAEModel against the official
daVinci-MagiHuman SAAudioFeatureExtractor.decode() path. The sibling
test_magi_human_sa_audio_parity.py validates against Diffusers
AutoencoderOobleck; this test validates against the upstream repo's custom
integration layer, which rebuilds AudioAutoencoder from model.pretransform.config
and filters pretransform.model.* weights.

Test passes at machine-eps (diff_max=0, diff_mean=0) in fp32, confirming
FV's SAAudioVAEModel is bit-identical to the official decode path. This
rules out the audio VAE as a contributor to OQ-6 (pipeline compounding).

Requires the daVinci-MagiHuman upstream clone and the gated
stabilityai/stable-audio-open-1.0 repo. Skips cleanly when either is absent.
2026-05-05 10:55:49 -07:00
SolitaryThinker d5acc7bbae [fix] magi-human: complete neg prompt; raise on missing CFG embeds
Complete the magi-human negative prompt to match upstream's three-block
concatenation (video + audio-quality + speech-delivery negatives). FV's
preset previously included only the video-side block, leaving audio CFG
to amplify the missing-block delta 5x via `v = uncond + 5 * (cond - uncond)`.

This explains the audio-side amplification observed in the pipeline-trace
investigation (step 1 v_cond_audio: 4.6% drift; v_cfg_audio: 13.8%).
Upstream reference: daVinci-MagiHuman/inference/pipeline/video_generate.py
:222-224.

Also harden CFG path: missing negative embeds previously silently fell
back to zeros, which is a hidden CFG amplifier. Now raises ValueError
explicitly when CFG=2 is set without negative embeds. Tracks OQ-6
production fix in tests/local_tests/magi-human.md.
2026-05-05 10:55:49 -07:00
SolitaryThinker 9035927da1 [docs] magi-human: numerical-alignment investigation findings + OQ-6/OQ-7
Document Wave 1-5 numerical-alignment investigation results in
tests/local_tests/magi-human.md (417 lines total):

- OQ-6: DiT bf16 noise-floor drift (diff_max=0.057 at atol=0.03);
  pipeline compounding at 4 steps (18.85x ratio vs 1-step baseline);
  root cause under investigation via _debug_magi_human_block_parity.py
- OQ-7: Wan VAE fp32 op-order drift (z*std+mean vs z/(1/std)+mean);
  shared Wan-family bug, deferred pending upstream fix

Also records Wave 1-5 methodology: loader fix, per-side log emission,
tolerance tightening, scheduler split, and bisect confirmation that
the compounding bug predates Wave 1.
2026-05-05 10:55:49 -07:00
SolitaryThinker 42d1c79694 [test] magi-human: surface bf16 bugs (OQ-6/OQ-7); 4-step pipeline; defer Wan VAE
Surface known bf16 numerical alignment bugs by tightening parity bounds:
- DiT parity atol=0.1 -> atol=0.03 (test FAILs at observed diff_max=0.057;
  bf16-noise-floor; tracked as OQ-6 root-cause investigation)
- Pipeline parity num_inference_steps=1 -> 4 (test FAILs at video
  diff_max=15.10, diff_mean=1.30 = 18.85x compounding ratio vs 1-step
  baseline 0.069; pre-existing compounding bug in original magi port,
  NOT introduced by Wave 1 -- bisect confirmed; tracked as OQ-6)
- Wan VAE parity atol=5e-2 -> atol=1e-3 (Wan VAE shared fp32 op-order
  drift, FV `z * std + mean` vs upstream `z / (1/std) + mean`; affects
  all Wan-family pipelines, deferred per OQ-7)

Also split _build_schedulers into _build_fastvideo_schedulers (double-shift,
matching MagiHumanDenoisingStage production) and _build_upstream_schedulers
(single-shift, matching MagiEvaluator.eval_with_text) to faithfully mirror
each side's production scheduler init pattern.

These tests SHIP IN A FAILING STATE intentionally as a forcing function
for follow-up investigation. See OQ-6 + OQ-7 in magi-human.md.
2026-05-05 10:55:49 -07:00
SolitaryThinker 59000cb933 [test] magi-human: emit per-side layer log files in block parity debugger
Add _debug_magi_human_block_parity.py, a standalone script that runs
forward hooks on both the upstream DiTModel and FastVideo MagiHumanDiT,
then writes per-side layer activation logs to:
  /tmp/opencode/magi_dit_up_layers.log
  /tmp/opencode/magi_dit_fv_layers.log

The logs record (name, shape, abs_mean, sum, min, max) for every hooked
activation, sorted by layer order. This enables side-by-side diff via
`diff /tmp/opencode/magi_dit_up_layers.log /tmp/opencode/magi_dit_fv_layers.log`
to pinpoint the first block where numerical divergence appears — a key
diagnostic step for OQ-6 root-cause investigation.
2026-05-05 10:55:49 -07:00
SolitaryThinker e6b15fc4de [test] magi-human: fix _find_base_shard_dir + snapshot_download fallback
Replace hf_hub_download (index-only canary) with snapshot_download using
allow_patterns=["base/*.safetensors", "base/model.safetensors.index.json"].
The old approach returned the parent of the index file, which could be a
symlink-resolved HF cache path that lacked the actual shard blobs. The new
approach verifies that at least one .safetensors shard is present before
returning the candidate directory, preventing false-positive skip decisions.

Applies to test_magi_human_parity.py, test_magi_human_pipeline_parity.py,
and the new _debug_magi_human_weight_diff.py debug script (which uses the
same loader pattern for weight-diff analysis).
2026-05-05 10:55:49 -07:00
SolitaryThinker 818daea816 [docs]: add tests/local_tests/magi-human.md (Phase 11 status) 2026-05-05 10:55:49 -07:00
SolitaryThinker d1bd1d8da4 [test]: principled bf16+CFG bounds for MagiHuman pipeline parity 2026-05-05 10:55:49 -07:00
SolitaryThinker d210a076e6 [misc]: name upstream constants in MagiHuman audio decoding 2026-05-05 10:55:49 -07:00
SolitaryThinker 7e96218d04 [perf]: cache RoPE+masks+kv_repeat in MagiHuman DiT forward 2026-05-05 10:55:49 -07:00
SolitaryThinker 6d0ddb44bc [perf]: hoist invariants out of MagiHuman denoising/latent-prep 2026-05-05 10:55:49 -07:00
SolitaryThinker 14c142ba9c [bugfix]: drop double-shift in MagiHuman scheduler init 2026-05-05 10:55:49 -07:00
SolitaryThinker 3fa3f4e333 [refactor] magi-human: defer audio VAE to main's stable-audio infra
Drop the magi-introduced 'encoder-style' Stable Audio VAE wrapper now
that main has merged a first-class Oobleck VAE port plus a shared
SAAudioVAEModel pipeline-glue lazy-loader (#1260). MagiHuman now
reuses that infrastructure instead of carrying duplicates:

  - audio VAE config:  fastvideo.configs.models.vaes.OobleckVAEConfig
  - audio VAE wrapper: fastvideo.models.vaes.sa_audio.SAAudioVAEModel

Removes:
  - fastvideo/configs/models/encoders/sa_audio.py (SAAudioVAEConfig)
  - fastvideo/models/encoders/sa_audio.py        (SAAudioVAEModel dup)

Updates:
  - magi-human pipeline + config to construct OobleckVAEConfig and set
    pretrained_path instead of arch_config.sa_audio_model_path.
  - encoders/__init__.py drops the SAAudioVAEConfig export.
  - parity test moves from tests/local_tests/encoders/ to .../vaes/
    (matches main's classification) and switches to the shared wrapper;
    explicit pretrained_dtype='float32' override keeps the fp32 ref
    parity check intact (default is fp16 to match official
    stable_audio_tools).

audio_decoding.py docstring still mentions 'sa_audio_vae_model' — main's
wrapper exposes that name as a back-compat alias for oobleck_vae, so
the existing comment is still accurate.
2026-05-05 10:55:49 -07:00
SolitaryThinker 5c7cd391ac [refactor] magi: T2V→base rename + faithful parity tests
- Drop misleading "T2V" framing — base MagiHuman is a joint audio-visual
  generator. Rename MagiHumanT2VConfig→MagiHumanBaseConfig,
  MagiHumanDistillT2VConfig→MagiHumanDistillConfig, presets
  magi_human_base_t2v→magi_human_base, magi_human_distill_t2v→
  magi_human_distill. Update comments / docstrings / example output
  filename. Keep workload_type="t2v" string (framework enum has no
  T2AV variant yet — same placeholder Stable Audio uses for T2A).
- Pipeline parity test: tighten atol/rtol from 0.5/0.5 to 0.35/0.05.
  atol absorbs the observed worst-element drift (~0.31 on signal abs_mean
  ~2.4); the tight rtol still flags gross structural bugs.
- Wan VAE parity: import FastVideo's AutoencoderKLWan from
  fastvideo.models.vaes.wanvae instead of diffusers. The test now
  actually validates the class MagiHumanBaseConfig.vae_config resolves
  to in production.
- Pipeline parity scheduler init: split per-side so upstream mirrors
  MagiEvaluator's single-shift pattern (fresh FlowUniPCMultistepScheduler
  + shift only in set_timesteps) and FastVideo mirrors
  MagiHumanDenoisingStage's double-shift pattern (constructor + set_timesteps
  both with shift). Surfaces a real production divergence between
  FastVideo's magi denoise loop and the official one for N>1 steps.
2026-05-05 10:55:49 -07:00
SolitaryThinker e7d6c40860 [feat] Port daVinci-MagiHuman base model with AV output
Ports GAIR-NLP/daVinci-MagiHuman's 15B-param joint-AV DiT into FastVideo:
text + video + audio joint denoise in one flat token stream, 32-step
FlowUniPC w/ CFG=2, Wan 2.2 TI2V-5B video VAE, Stable Audio Open 1.0
audio VAE (first-class port at fastvideo/models/vaes/oobleck.py, no
runtime diffusers import), T5-Gemma 9B UL2 text encoder, auto-muxed
mp4 with h264 video + aac stereo audio.

Parity tests all pass on converted weights: DiT (diff_median=2.6e-3),
text encoder (exact), video VAE (8e-4), audio VAE (exact), pipeline
latent (5e-2).

See .agents/skills/add-model/REVIEW.md for porting procedure + open
review items.
2026-05-05 10:55:49 -07:00
SolitaryThinker 12a7cb53ff skill draft 2026-05-05 10:55:49 -07:00
MookandSolitaryThinker c17d33bf33 [ci] Replace flaky LTX-2 pixel SSIM with latent-slice cosine regression (#1253)
Co-authored-by: SolitaryThinker <wlsaidhi@gmail.com>
2026-05-05 03:36:58 -07:00
2aaeee2ab8 [feat] Improve API: streaming router (multi-replica load balancer + ws proxy) (#1286)
Co-authored-by: Junda (David) Su <90978028+Davids048@users.noreply.github.com>
Co-authored-by: Matthew Noto <99706358+RandNMR73@users.noreply.github.com>
Co-authored-by: XOR-op <17672363+XOR-op@users.noreply.github.com>
Co-authored-by: Zhang Peiyuan <42993249+jzhang38@users.noreply.github.com>
2026-05-05 03:00:08 -07:00
eb3a394224 [feat] Improve API: streaming auxiliaries (safety, rewrite, logger, mock) (#1284)
Co-authored-by: Junda (David) Su <90978028+Davids048@users.noreply.github.com>
Co-authored-by: Matthew Noto <99706358+RandNMR73@users.noreply.github.com>
Co-authored-by: XOR-op <17672363+XOR-op@users.noreply.github.com>
Co-authored-by: Zhang Peiyuan <42993249+jzhang38@users.noreply.github.com>
2026-05-05 00:14:34 -07:00
f673423b51 [feat] Improve API: streaming prompt enhancer with LLMProvider abstraction (#1258)
Co-authored-by: Junda (David) Su <90978028+Davids048@users.noreply.github.com>
Co-authored-by: Matthew Noto <99706358+RandNMR73@users.noreply.github.com>
Co-authored-by: XOR-op <17672363+XOR-op@users.noreply.github.com>
Co-authored-by: Zhang Peiyuan <42993249+jzhang38@users.noreply.github.com>
2026-05-04 13:44:40 -07:00
eb0a41528a [feat] Improve API: streaming server GpuPool + worker subprocess (#1257)
Co-authored-by: Junda (David) Su <90978028+Davids048@users.noreply.github.com>
Co-authored-by: Matthew Noto <99706358+RandNMR73@users.noreply.github.com>
Co-authored-by: XOR-op <17672363+XOR-op@users.noreply.github.com>
Co-authored-by: Zhang Peiyuan <42993249+jzhang38@users.noreply.github.com>
2026-05-04 12:56:31 -07:00
William Lin 140bd1a6cf [misc]: standardize install instructions on uv pip install (#1279) 2026-05-02 12:45:50 -07:00
William Lin 11f5a8e582 [misc] pin torch to 2.11.0 (#1277) 2026-05-02 11:48:07 -07:00
71b3cb8c34 [ci] Add CI Performance Regression Tracking Changes (#1248)
Co-authored-by: Satyam Srivastava <satyam53@Mac.lan1>
Co-authored-by: Satyam Srivastava <satyam53@Satyams-MacBook-Air.local>
Co-authored-by: SolitaryThinker <wlsaidhi@gmail.com>
2026-05-02 03:29:53 -07:00
William Lin c85f6a477f [docs] add hierarchical AGENTS.md per-directory guidance (#1278) 2026-05-02 03:28:18 -07:00
Junda Su 40d4930d73 [bugfix] Update fa import (#1271) 2026-05-02 01:25:22 -07:00
William Lin f9be085243 [ci] pre-commit: drop stale excludes + document agent lint flow (#1276) 2026-05-02 01:19:06 -07:00
William Lin 36b53ff350 [bugfix]: classify stable_audio fields in schema parity inventory (#1275) 2026-05-02 00:12:10 -07:00
William Lin 9801037c3d [refactor] tests/local_tests: organize by model family (#1269) 2026-05-01 01:49:54 -07:00
alexzmsandmergify[bot] 74d09b0efd [misc] cleanup: grad-norm asserts, dead offload file, callback names (#1268)
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
2026-05-01 01:16:13 -07:00
alexzms 38dc8820ac [ci] add CPU unit tests for train callback system in fastvideo.train (#1267) 2026-05-01 00:53:29 -07:00
William Lin c77a76c6af [feat] Stable Audio Open 1.0: T2A + A2A + RePaint inpainting (native) (#1260) 2026-05-01 00:07:11 -07:00
alexzms d14d5aadea [feat] Cosmos 2.5 training support in fastvideo.train (#1224) 2026-05-01 01:15:02 +00:00
alexzms 4c915b7742 [ci] add CPU unit tests for train checkpoint utilities in fastvideo.train (#1265) 2026-04-29 18:55:39 +00:00
alexzms 9a8bbe18fa [bugfix]: fix SP deadlock in negative prompt encoding during training (#1178) 2026-04-28 01:06:49 +00:00
alexzms ea25441ef0 [ci] add CPU unit tests for fastvideo.train load_run_config (#1264) 2026-04-28 01:06:18 +00:00
48957fcde1 [bugfix] Fix modal remote functions crash container on sys exit in CI remote functions (#1261)
Co-authored-by: Satyam Srivastava <satyam53@Mac.lan1>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
2026-04-27 21:50:55 +00:00
Mook 7b872cc41e [Perf] Skip bool-mask round-trip in block-sparse VSA attention (#1243) 2026-04-26 15:14:37 -07:00
alexzms 37418946c8 [docs]: clarify real_score_guidance_scale CFG parameterization (#1256) 2026-04-26 16:38:00 +08:00
William Lin 95fd29e0cb [feat] Streaming WebSocket server skeleton (single generator + fMP4) (#1251) 2026-04-26 00:33:49 -07:00
Junda Suandmergify[bot] e17cd2633c [bugfix]: normalize uint8 pil_image in I2V VAE encoding (#1249)
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
2026-04-24 09:16:01 +00:00
William Lin e0dc5f2b0c [feat] Add typed LTX-2 continuation state and streaming session store (#1250) 2026-04-24 01:28:07 -07:00
William Lin 70ee5d230c [feat] [6/n] Improve API: LTX-2 public preset + asset wiring + gpu_pool translation (#1239) 2026-04-23 11:36:45 -07:00
William Lin 24ced500f5 [test] add LTX-2 distilled T2V SSIM regression test (#1240) 2026-04-21 12:03:38 -07:00
298 changed files with 29710 additions and 987 deletions
+1 -1
View File
@@ -119,7 +119,7 @@ FastVideo-WorldModel/
## Build & Test Commands
```bash
uv pip install -e .[dev] # Editable install
uv pip install -e ".[dev]" # Editable install
pre-commit run --all-files # Lint/format/spell
pytest tests/ # Top-level tests
pytest fastvideo/tests/ -v # Package tests
+96
View File
@@ -0,0 +1,96 @@
#!/usr/bin/env bash
# Sync .agents/skills/ into .claude/skills/ via per-skill symlinks.
#
# Why: Claude Code only scans .claude/skills/ and ~/.claude/skills/ for
# user-invocable skills (no skillsPath config exists — see
# https://code.claude.com/docs/en/skills.md). This repo's skills live
# in .agents/skills/ so they travel with the repo and stay under git.
# Run this once after cloning (or after adding/removing a skill) to
# expose them to Claude Code without maintaining a parallel tree.
#
# Usage:
# .agents/scripts/sync-skills.sh
#
# Idempotent and safe to re-run. Prunes stale symlinks whose source
# has been removed from .agents/skills/. Leaves hand-written
# .claude/skills/<name>/ directories untouched (only symlinks are
# managed).
set -euo pipefail
REPO_ROOT="$(git -C "$(dirname "$0")" rev-parse --show-toplevel)"
SRC_DIR="$REPO_ROOT/.agents/skills"
DST_DIR="$REPO_ROOT/.claude/skills"
if [[ ! -d "$SRC_DIR" ]]; then
echo "Error: $SRC_DIR does not exist." >&2
exit 1
fi
mkdir -p "$DST_DIR"
linked=0
unchanged=0
skipped=0
pruned=0
link_skill() {
local name="$1"
local src="$SRC_DIR/$name"
local dst="$DST_DIR/$name"
# Relative target keeps symlinks portable across clones.
local rel="../../.agents/skills/$name"
if [[ -L "$dst" ]]; then
if [[ "$(readlink "$dst")" == "$rel" ]]; then
unchanged=$((unchanged + 1))
return
fi
rm "$dst"
elif [[ -e "$dst" ]]; then
echo "Skipped (not a symlink): .claude/skills/$name" >&2
skipped=$((skipped + 1))
return
fi
ln -s "$rel" "$dst"
echo "Linked: .claude/skills/$name -> $rel"
linked=$((linked + 1))
}
prune_stale() {
local link="$1"
local target
target="$(readlink "$link")"
case "$target" in
../../.agents/skills/*) ;;
*) return ;;
esac
local name="${target##*/}"
if [[ ! -d "$SRC_DIR/$name" ]]; then
rm "$link"
echo "Pruned stale: .claude/skills/$(basename "$link")"
pruned=$((pruned + 1))
fi
}
for src in "$SRC_DIR"/*/; do
[[ -d "$src" ]] || continue
name="$(basename "$src")"
# Only treat directories that actually contain a SKILL.md as skills.
[[ -f "$src/SKILL.md" ]] || continue
link_skill "$name"
done
shopt -s nullglob
for link in "$DST_DIR"/*; do
[[ -L "$link" ]] || continue
prune_stale "$link"
done
shopt -u nullglob
printf "\nSummary: %d linked, %d unchanged, %d pruned" "$linked" "$unchanged" "$pruned"
if [[ "$skipped" -gt 0 ]]; then
printf ", %d skipped (non-symlink collision)" "$skipped"
fi
printf "\n"
+3
View File
@@ -5,3 +5,6 @@
{"name": "evaluate-video-quality", "description": "Evaluate generated video quality using available metrics (SSIM, loss trajectory, caption consistency)", "path": "evaluate-video-quality/SKILL.md", "status": "draft", "trust": "low"}
{"name": "index-related-work", "description": "Ingest a paper or repository into the related work index", "path": "index-related-work/SKILL.md", "status": "draft", "trust": "low"}
{"name": "search-related-work", "description": "Query the related work index for relevant papers, repos, or comparisons", "path": "search-related-work/SKILL.md", "status": "draft", "trust": "low"}
{"name": "seed-ssim-references", "description": "Run a new or updated fastvideo/tests/ssim/ test on Modal, pull generated videos, and upload them to FastVideo/ssim-reference-videos so the test has a regression baseline", "path": "seed-ssim-references/SKILL.md", "status": "draft", "trust": "low"}
{"name": "reseed-ssim-references", "description": "Re-seed (overwrite) HF reference videos for an existing fastvideo/tests/ssim/ test and a single model id on Modal L40S. Always backs up current refs first, regenerates on Modal, pauses for the user to eyeball before-vs-after, then uploads with --force scoped to --model-id. Sister skill to seed-ssim-references; use when intentional code change has invalidated existing refs", "path": "reseed-ssim-references/SKILL.md", "status": "draft", "trust": "low"}
{"name": "add-model", "description": "Add a new model (or variant) to FastVideo: DiT + configs + pipeline + presets + registry + tests. Walks through FastVideo's single stage-based pipeline architecture with exact file paths and registration hooks.", "path": "add-model/SKILL.md", "status": "draft", "trust": "low"}
+1 -1
View File
@@ -12,7 +12,7 @@ automates the boilerplate of setting environment variables, picking the right
entrypoint, and applying defaults from the closest example script.
## Prerequisites
- The repo is cloned and `fastvideo` is installed (`uv pip install -e .[dev]`).
- The repo is cloned and `fastvideo` is installed (`uv pip install -e ".[dev]"`).
- Dataset is preprocessed (see `docs/training/data_preprocess.md`).
- `WANDB_API_KEY` is set in the environment (or `WANDB_MODE=offline` for local).
- GPU resources are available (multi-GPU requires NCCL).
@@ -0,0 +1,343 @@
---
name: reseed-ssim-references
description: Re-seed HF reference videos for a single existing SSIM test on Modal L40S. Always backs up current refs locally first, regenerates on Modal, pauses for the user to eyeball before-vs-after quality, then overwrites the targeted `<model_id>` subtree on `FastVideo/ssim-reference-videos` with `--force`. Use when an intentional code change (model port fix, attention backend swap, kernel upgrade, hyperparameter change) has invalidated existing refs and they need to be regenerated. Pairs with `seed-ssim-references`, which is for first-time seeding only.
---
# Re-seed SSIM Reference Videos
## Purpose
Replace the existing SSIM reference videos for a single `(test_file, model_id)`
pair on the HF dataset (`FastVideo/ssim-reference-videos`). This is **destructive**
on HF — the old refs are overwritten — so the skill always:
1. Confirms intent with a one-liner the user has to type.
2. Downloads the existing refs as a local, timestamped backup.
3. Regenerates on Modal L40S (same code path that CI uses).
4. Pauses for a side-by-side eyeball of backup vs new mp4s.
5. Uploads with `--force`, scoped to the single `--model-id`.
6. Reminds the user to keep the backup until the PR lands.
Pairs with `seed-ssim-references`, which is the inverse (first-time seeding
only, refuses to overwrite). Re-seeding is intentionally a separate, more
ceremonial operation because mistakenly clobbering production refs is much
harder to recover from than failing closed.
## When to use
- An intentional code change (model port fix, kernel upgrade, attention
backend swap, hyperparameter change in the test itself) has shifted the
expected SSIM output and the existing refs no longer represent the new
ground truth.
- A test is failing in CI **for the right reason** (the new code is correct,
the old refs are stale).
## When not to use
- A test is failing for the **wrong** reason (the port is buggy, not the
refs). Fix the port; re-seeding hides the bug.
- A brand-new test that has no refs on HF yet. Use `seed-ssim-references`.
- "Just to clean up drift" without a concrete code change to point at. The
PR description has to justify *why* refs changed; without a concrete
change, there's nothing to write.
## Inputs
| Parameter | Required | Description |
|-----------|----------|-------------|
| `test_file` | Yes | Path to the SSIM test, e.g. `fastvideo/tests/ssim/test_matrixgame_similarity.py`. Validated against `fastvideo/tests/ssim/test_*_similarity.py`. |
| `model_id` | Yes | Single model id from the test's `*_MODEL_TO_PARAMS`, e.g. `Matrix-Game-2.0-Diffusers-Base`. Re-seed runs are **per model**. For multi-model tests, invoke the skill once per model. |
| `intent_rationale` | Yes | One-line explanation of *why* refs are being regenerated (e.g. "Relax FA-2 head_size whitelist to include 80 — matrix_game now uses FLASH_ATTN instead of TORCH_SDPA"). Recorded in the backup directory and reused in the PR description. |
Hardcoded:
- Modal GPU: **L40S** (matches CI; re-seeding from another SKU produces refs
that L40S CI cannot match).
- Quality tier: **`default`**. `full_quality` is a separate, deliberate
operation.
- HF repo: `FastVideo/ssim-reference-videos` (override via
`FASTVIDEO_SSIM_REFERENCE_HF_REPO`).
- Device folder: `L40S_reference_videos`.
## Prerequisites
The user has confirmed:
- `modal` CLI authenticated.
- `hf` CLI authenticated, **and** `HF_API_KEY` (or `HUGGINGFACE_HUB_TOKEN` /
`HF_TOKEN`) exported with **write** access to
`FastVideo/ssim-reference-videos`.
- The current branch's code is the change that motivated the re-seed (i.e.
`git rev-parse HEAD` is the commit that intentionally invalidated refs).
Fail fast if any of these are missing.
## Steps
### 1. Validate inputs and confirm intent
- Verify `test_file` exists and matches `fastvideo/tests/ssim/test_*_similarity.py`.
- Grep the file for `*_MODEL_TO_PARAMS` and assert `model_id` is one of its
keys. If the file has only a single hardcoded model, accept that model id
as the only valid value.
- Print the rationale and ask the user to type **`confirm reseed`** (not just
`y` — make it deliberate):
> About to RE-SEED references for model `<model_id>` from test `<test_file>`.
> This will OVERWRITE existing refs on
> `FastVideo/ssim-reference-videos/reference_videos/default/L40S_reference_videos/<model_id>/`
> after backup + Modal regen + eyeball.
>
> Reason: `<intent_rationale>`
> HEAD: `<git rev-parse --short=12 HEAD>`
>
> Reply `confirm reseed` to proceed, anything else to abort.
Stop until the user types exactly `confirm reseed`. Anything else aborts
with no side effects.
### 2. Back up existing refs
Always required. The backup is the only graceful path back if anything goes
wrong later.
```bash
SHORT_COMMIT=$(git rev-parse --short=12 HEAD)
TIMESTAMP=$(date -u +%Y%m%d_%H%M%S)
MODEL_SAFE=$(echo "<model_id>" | tr '/' '_')
BACKUP_DIR="ssim_reseed_backup/${TIMESTAMP}_${SHORT_COMMIT}_${MODEL_SAFE}"
mkdir -p "$BACKUP_DIR"
hf download \
--repo-type dataset FastVideo/ssim-reference-videos \
--include "reference_videos/default/L40S_reference_videos/<model_id>/**" \
--local-dir "$BACKUP_DIR"
mp4_count=$(find "$BACKUP_DIR" -name "*.mp4" | wc -l)
echo "Backup mp4 count: $mp4_count"
[ "$mp4_count" -gt 0 ] || {
echo "ERROR: backup is empty for <model_id>. Either the model id is wrong"
echo "or there are no existing refs (use seed-ssim-references instead)."
exit 1
}
# Provenance — used in the PR description
cat > "$BACKUP_DIR/PROVENANCE.txt" <<EOF
test_file: <test_file>
model_id: <model_id>
head_commit: $(git rev-parse HEAD)
timestamp_utc: $(date -u +%FT%TZ)
reason: <intent_rationale>
EOF
```
If the `hf download` produces zero mp4s, abort — the user has either picked a
non-existent `model_id` or there are no refs yet (in which case
`seed-ssim-references` is the right tool).
### 3. Regenerate on Modal L40S
Mirror CI's exact env recipe so the regenerated refs are byte-comparable to
what CI will produce on the same commit. Two differences from CI:
1. **Pass the same env prefix CI uses** (`IMAGE_VERSION`, `BUILDKITE_*`) — see
`.buildkite/pipeline.yml:1-3` and `.buildkite/scripts/pr_test.sh:62-83`.
Without this, `ssim_test.py:17-18` resolves a different GHCR image tag
(default is `latest`, CI is `py3.12-latest`), and `ssim_test.py:38-46`
bakes different values into the image's frozen env block. **Mismatched
image or env is the most common source of SSIM drift between reseed and
CI runs.**
2. **Do not pass `--skip-reference-download`**. Letting the test fetch the
existing refs and run the full SSIM compare gives "before" SSIM numbers
for the PR description, and the test still produces the new mp4s
regardless of whether the comparison passes or fails.
```bash
SUBDIR="${TIMESTAMP}_${SHORT_COMMIT}"
IMAGE_VERSION="py3.12-latest" \
BUILDKITE_REPO="$(git config --get remote.origin.url)" \
BUILDKITE_COMMIT="$(git rev-parse HEAD)" \
BUILDKITE_PULL_REQUEST="${BUILDKITE_PULL_REQUEST:-false}" \
modal run fastvideo/tests/modal/ssim_test.py \
--git-repo="$(git config --get remote.origin.url)" \
--git-commit="$(git rev-parse HEAD)" \
--hf-api-key="$HF_API_KEY" \
--test-files="<test_file>" \
--sync-generated-to-volume \
--generated-volume-subdir="$SUBDIR" \
--no-fail-fast
```
Capture the printed `modal volume get ...` hint — its `<SUBDIR>` matches
`$SUBDIR` and is needed for step 4. Capture the SSIM numbers from the test
output (or from the JSON next to the generated mp4) for the PR description.
### 4. Download generated videos
```bash
modal volume get --force hf-model-weights \
ssim_generated_videos/default/"$SUBDIR"/generated_videos \
./generated_videos_modal/default
```
After this, the new mp4s live at:
```
./generated_videos_modal/default/generated_videos/L40S_reference_videos/<model_id>/<backend>/<prompt>.mp4
```
`--force` is required when `./generated_videos_modal/default` already exists
from a prior run; safe on the first run too.
### 5. PAUSE — user reviews quality side-by-side
Print the diff and the comparison:
```bash
echo "=== File list diff (backup vs new) ==="
diff -u \
<(find "$BACKUP_DIR/reference_videos/default/L40S_reference_videos/<model_id>" -name "*.mp4" \
| sed "s|$BACKUP_DIR/reference_videos/default/L40S_reference_videos/||" | sort) \
<(find ./generated_videos_modal/default/generated_videos/L40S_reference_videos/<model_id> -name "*.mp4" \
| sed "s|./generated_videos_modal/default/generated_videos/L40S_reference_videos/||" | sort) \
|| true
echo
echo "=== SSIM numbers from this run (paste into PR) ==="
find ./generated_videos_modal/default/generated_videos/L40S_reference_videos/<model_id> -name "*_ssim.json" -exec cat {} \;
```
Then stop and tell the user:
> Old refs backed up to `$BACKUP_DIR`.
> New videos in `./generated_videos_modal/default/generated_videos/L40S_reference_videos/<model_id>/`.
>
> Open both in a video player. Confirm the new videos:
> 1. Look correct (no obvious artifacts, no black/static frames).
> 2. Are *intentionally* different from the backup in the way described
> in `<intent_rationale>` (e.g. slight numerical drift only, not a
> different scene / different motion / corrupted output).
>
> Reply **`upload`** to overwrite HF, anything else to abort.
> Aborting leaves the backup and new videos on disk for inspection — nothing
> on HF changes.
Do not proceed until the user types exactly `upload`. If they abort, leave
everything on disk and stop here.
### 6. Copy into the local reference layout
Same as `seed-ssim-references` step 5:
```bash
python fastvideo/tests/ssim/reference_videos_cli.py copy-local \
--quality-tier default \
--device-folder L40S_reference_videos \
--generated-dir ./generated_videos_modal/default/generated_videos/L40S_reference_videos
```
Result: `fastvideo/tests/ssim/reference_videos/default/L40S_reference_videos/<model_id>/<backend>/<prompt>.mp4`.
### 7. Upload with `--force`, scoped to `--model-id`
The `--force` flag is what makes this skill different from `seed-ssim-references`.
Always pair it with `--model-id` so a typo cannot accidentally overwrite a
neighboring model's refs.
```bash
python fastvideo/tests/ssim/reference_videos_cli.py upload \
--quality-tier default \
--device-folder L40S_reference_videos \
--model-id "<model_id>" \
--force
```
The CLI's overwrite guard refuses without `--force`; with `--force` it
overwrites only files under
`reference_videos/default/L40S_reference_videos/<model_id>/`.
### 8. Report success and retention guidance
Print:
- The HF path that was overwritten (`<repo>/reference_videos/default/L40S_reference_videos/<model_id>/`).
- The local backup directory path.
- The new SSIM numbers from step 5.
- This restore command, in case the PR review surfaces a problem after
upload:
```bash
python fastvideo/tests/ssim/reference_videos_cli.py upload \
--quality-tier default \
--device-folder L40S_reference_videos \
--model-id "<model_id>" \
--reference-dir "$BACKUP_DIR/reference_videos/default/L40S_reference_videos" \
--force
```
- This PR-description checklist (see `fastvideo/tests/ssim/AGENTS.md` →
*Updating Reference Videos*):
1. Source commit that produced the new refs (HEAD at re-seed time).
2. Test command and GPU SKU (`L40S`).
3. Before/after SSIM numbers.
4. The `<intent_rationale>` from step 1.
5. A note that the backup lives at `$BACKUP_DIR` and should be retained
until CI on the PR is green.
Do **not** auto-rerun the SSIM test — the user does that as part of the PR.
## Failure modes and how to handle them
- **`HF_API_KEY` unset.** Stop before step 2.
- **Backup is empty (zero mp4s).** Stop before step 3 — the model id is
wrong or the refs don't exist yet (use `seed-ssim-references`).
- **Modal run fails before generation.** No mp4s on the volume. Don't
upload. Investigate the failure (test crash, OOM, partition exhaustion),
fix, then retry from step 3. Backup is still intact.
- **Quality regressed (visual or metric).** User aborts at step 5. Backup
retained. New videos retained on disk for inspection. Nothing on HF
changed. Either fix the underlying code change or abandon the re-seed.
- **User confirmed `upload` but later realized the new refs are wrong.**
Run the restore command from step 8 with the backup `--reference-dir`.
This is exactly why the backup exists.
- **Multi-model test, only one model is being re-seeded.** Run the skill
once per model id. The `--model-id` scope on upload guarantees the others
are untouched.
## Design notes (for future skill maintainers)
- Per-`model_id` scope is mandatory. The dataset houses many model subtrees;
re-seeding the wrong one is hard to undo without backup.
- `default` tier only; `full_quality` is a separate, deliberate operation
with different params and ~doubled runtime, and isn't what CI gates on.
- The skill deliberately does **not** pass `--skip-reference-download` to
Modal so we get pre-reseed SSIM numbers for the PR. The `seed`-skill
passes it because no refs exist yet; for re-seed, refs do exist and
exposing the comparison is informative.
- The two-token confirm (`confirm reseed`, then `upload`) is intentional.
Re-seeding is high-blast-radius and should not be one-keystroke.
- The backup directory is plain mp4s + `PROVENANCE.txt`. No HF metadata is
preserved; the restore path uses `reference_videos_cli.py upload
--reference-dir` which doesn't need it.
## References
- `.agents/skills/seed-ssim-references/SKILL.md` — the first-time seed
skill this one parallels. Read it for the Modal flag rationale shared
between the two flows.
- `fastvideo/tests/ssim/AGENTS.md` — directory rules, including the PR
expectations for any reference-video change (rationale, before/after
SSIM, source commit/model/backend).
- `fastvideo/tests/ssim/reference_videos_cli.py` — `copy-local`, `upload`
(with `--model-id`, `--force`), `download`. The overwrite guard at
`upload_reference_videos` is the safety net this skill leans on.
- `fastvideo/tests/modal/ssim_test.py` — Modal orchestrator;
`--sync-generated-to-volume`, `--generated-volume-subdir`,
`--skip-reference-download`, `--no-fail-fast`.
## Changelog
| Date | Change |
|------|--------|
| 2026-05-02 | Initial version. Sister skill to `seed-ssim-references`, scoped to single `(test_file, model_id)` re-seeds, with mandatory backup and two-token confirm. |
@@ -0,0 +1,376 @@
---
name: seed-ssim-references
description: Seed HF reference artefacts for a single newly-added SSIM test (pixel `.mp4` for `run_text_to_video_similarity_test`-style tests, or latent `.pt` for `run_text_to_latent_similarity_test`-style tests). Runs the test on Modal L40S, downloads the generated artefacts via `modal volume get`, pauses for the user to verify (visual eyeball for mp4, numerics dump for pt), then uploads only that test's files to `FastVideo/ssim-reference-videos`. Use when a new `fastvideo/tests/ssim/test_*_similarity.py` has just been added and has no references on HF yet.
---
# Seed SSIM Reference Artefacts (mp4 or pt)
## Purpose
A brand-new SSIM test in `fastvideo/tests/ssim/` fails forever until its
reference artefacts exist on the HF dataset
(`FastVideo/ssim-reference-videos`). The dataset hosts two kinds of artefacts
side-by-side per `(model_id, backend, prompt)`:
- **`.mp4`** — pixel ground-truth for tests that call
`run_text_to_video_similarity_test` / `run_image_to_video_similarity_test`
in `inference_similarity_utils.py`. Compared via SSIM.
- **`.pt`** — pre-VAE latent bundle (fp16 full latent + fp32 slice +
metadata + `slice_spec` + `format_version`) for tests that call
`run_text_to_latent_similarity_test` in `latent_similarity_utils.py`.
Compared via cosine distance on the slice and the full tensor.
This skill:
1. Detects which artefact type the test produces (pixel vs latent).
2. Runs the test on Modal's L40S pool to generate the artefacts.
3. Downloads them to the local repo via `modal volume get`.
4. Pauses so the user can verify quality:
- **mp4**: visual eyeball in a video player.
- **pt**: numerics dump (shape, slice stats, NaN/Inf check, metadata).
5. Uploads only the new test's files to HF, with a guard that refuses to
overwrite anything already present.
The skill is run **manually**, once per new test. Before invoking it, the user
has already sanity-tested the new test locally — it launches `VideoGenerator`
and writes an artefact without crashing (the missing-reference assertion at
the end is expected). The skill does not re-test locally; it goes straight
to Modal L40S (which is what CI uses).
## When to use
- A new `test_*_similarity.py` file has been added in `fastvideo/tests/ssim/`
and the HF dataset has no `reference_videos/default/L40S_reference_videos/<model_id>/`
subtree for it yet.
## When not to use
- Regular CI runs — once refs exist, `pytest fastvideo/tests/ssim/` downloads
them automatically.
- Re-seeding an existing test. That requires `--force` on the upload step, and
is out of scope here; treat as a separate, deliberate operation.
## Inputs
The skill has **one required input**: the path to the new SSIM test file.
Prompt the user for it if they didn't supply it.
| Parameter | Required | Description |
|-----------|----------|-------------|
| `test_file` | Yes | e.g. `fastvideo/tests/ssim/test_ltx2_similarity.py`. The skill's first action is to ask for this if missing. |
Everything else is fixed:
- Modal runner GPU: **L40S** (hardcoded in `fastvideo/tests/modal/ssim_test.py`).
- Device folder: `L40S_reference_videos`.
- Quality tier: `default` (the tier CI runs). The `full_quality` tier is not
seeded by this skill.
- HF repo: `FastVideo/ssim-reference-videos` (dataset).
- Multi-model test files: all model ids in `*_MODEL_TO_PARAMS` are seeded
together; the Modal run produces one mp4 per (model, prompt, backend) and
the upload scopes by `--model-id`, looping if there is more than one.
## Prerequisites
The user has confirmed:
- `modal` CLI authenticated.
- `HF_API_KEY` (or `HUGGINGFACE_HUB_TOKEN` / `HF_TOKEN`) exported with write
access to `FastVideo/ssim-reference-videos`.
- The test file runs locally end-to-end (generates an mp4; SSIM assertion
failure due to missing reference is expected and fine).
Fail fast if the token env var is missing.
## Steps
### 1. Ask for the test file, then detect artefact type
If the user didn't name one, ask: *"Which SSIM test file do you want to seed
references for? (e.g. `fastvideo/tests/ssim/test_ltx2_similarity.py`)"*.
Validate:
- Path exists and matches `fastvideo/tests/ssim/test_*_similarity.py`.
- File defines a `*_MODEL_TO_PARAMS` dict — grep it to extract the set of
model ids. Those ids drive step 5.
Detect artefact type by inspecting the file's imports / helper call:
- **latent** (`.pt`) — file imports `run_text_to_latent_similarity_test`
from `fastvideo.tests.ssim.latent_similarity_utils` (or any other helper
that ends with `_latent_similarity_test`).
- **pixel** (`.mp4`) — file imports
`run_text_to_video_similarity_test` / `run_image_to_video_similarity_test`
from `fastvideo.tests.ssim.inference_similarity_utils`, OR uses the
legacy custom-inline helper pattern (see `test_gamecraft`,
`test_longcat`, etc.). Default to pixel when both heuristics fail.
Record `ARTEFACT_TYPE ∈ {pixel, latent}` for use in step 4. Steps 2, 3, 5,
and 6 are artefact-type-agnostic — `_iter_reference_files`,
`copy_generated_to_reference`, and `upload_reference_videos` already walk
both `.mp4` and `.pt` (see `reference_videos_cli.py`).
If either check fails, stop and tell the user what's wrong.
### 2. Run the test on Modal L40S
Pick a subdir name so repeated runs don't collide:
```bash
SHORT_COMMIT=$(git rev-parse --short=12 HEAD)
TIMESTAMP=$(date -u +%Y%m%d_%H%M%S)
SUBDIR="${TIMESTAMP}_${SHORT_COMMIT}"
```
Then launch the Modal run. The `IMAGE_VERSION` and `BUILDKITE_*` env-prefix
**must** match what CI exports in `.buildkite/scripts/pr_test.sh`, otherwise
`fastvideo/tests/modal/ssim_test.py` resolves a different GHCR image tag
(default is `latest`, CI is `py3.12-latest`) and bakes different values into
the image's frozen env block (`ssim_test.py:17-18, 38-46`). Mismatched image
or env produces SSIM drift that doesn't show up until the same commit runs
in CI.
```bash
IMAGE_VERSION="py3.12-latest" \
BUILDKITE_REPO="$(git config --get remote.origin.url)" \
BUILDKITE_COMMIT="$(git rev-parse HEAD)" \
BUILDKITE_PULL_REQUEST="${BUILDKITE_PULL_REQUEST:-false}" \
modal run fastvideo/tests/modal/ssim_test.py \
--git-repo="$(git config --get remote.origin.url)" \
--git-commit="$(git rev-parse HEAD)" \
--hf-api-key="$HF_API_KEY" \
--test-files="<test_file>" \
--sync-generated-to-volume \
--generated-volume-subdir="$SUBDIR" \
--skip-reference-download \
--no-fail-fast
```
Env prefix rationale (parity with CI; see `.buildkite/pipeline.yml:1-3` and
`.buildkite/scripts/pr_test.sh:62-83`):
- `IMAGE_VERSION=py3.12-latest`: pins the Modal image tag to the same one CI
uses. Without this, `ssim_test.py:17` falls back to `latest`, which on
GHCR is built from `Dockerfile.python3.10` — different Python, torch, and
flash-attn wheel than CI's `py3.12-latest` (`infra-build-image.yml:51-67`,
`_template-build-image.yml:65-101`).
- `BUILDKITE_REPO`/`BUILDKITE_COMMIT`/`BUILDKITE_PULL_REQUEST`: mirror what
Buildkite exports. `ssim_test.py:38-46` bakes these into the image's
`.env(...)` block; mismatched values can perturb in-container code paths
that branch on PR-vs-non-PR. `false` for `BUILDKITE_PULL_REQUEST` matches
Buildkite's "non-PR build" sentinel.
Flag rationale:
- `--skip-reference-download`: no refs exist yet, so conftest must not try to
pull them.
- `--no-fail-fast`: lets the test finish generation before `_assert_similarity`
raises `FileNotFoundError: Reference video folder does not exist`. The
expected failure is what we want — the mp4 has already been written.
- `--sync-generated-to-volume` + `--generated-volume-subdir`: copies the
generated mp4s to the `hf-model-weights` Modal volume under
`ssim_generated_videos/default/<SUBDIR>/generated_videos/` so we can pull
them locally.
The Modal run will end with a nonzero exit (expected) and print a
`modal volume get hf-model-weights ssim_generated_videos/default/<SUBDIR>/generated_videos ./generated_videos_modal/default`
command. Capture that `<SUBDIR>` — you need it for step 3.
### 3. Download generated videos locally
```bash
modal volume get --force hf-model-weights \
ssim_generated_videos/default/"$SUBDIR"/generated_videos \
./generated_videos_modal/default
```
`--force` is required when the parent `./generated_videos_modal/default`
already exists; without it, `modal volume get` errors with `[Errno 21] Is a
directory`. Safe to pass on the first run too.
After this, the mp4s live at
`./generated_videos_modal/default/generated_videos/L40S_reference_videos/<model_id>/<backend>/<prompt>.mp4`.
The extra `generated_videos/` level comes from the volume layout in
`_sync_generated_videos_to_volume` (`ssim_test.py`) — the command copies
`<repo>/fastvideo/tests/ssim/generated_videos/<tier>` to
`ssim_generated_videos/<tier>/<SUBDIR>/generated_videos/`, and `modal volume
get` preserves that trailing `generated_videos/` segment.
### 4. PAUSE — user reviews quality
Type-aware verification.
**For `ARTEFACT_TYPE = pixel`** — list the downloaded mp4s and ask the user to
open them in a video player:
> "Generated videos downloaded to `./generated_videos_modal/default/generated_videos/L40S_reference_videos/`. Please open them and confirm the quality looks correct. Reply **`upload`** to continue, or anything else to abort."
**For `ARTEFACT_TYPE = latent`** — `.pt` files are not human-watchable. Print
a numerics dump for each `.pt` so the user can sanity-check shape, distribution,
and metadata:
```python
import torch
from pathlib import Path
ROOT = Path("./generated_videos_modal/default/generated_videos/L40S_reference_videos")
for p in sorted(ROOT.rglob("*.pt")):
d = torch.load(p, map_location="cpu", weights_only=False)
s = d["expected_slice"]
L = d["latent"].float()
print(f"=== {p.relative_to(ROOT)} ===")
print(f" format_version: {d['format_version']}")
print(f" shape: {d['shape']}")
print(f" dtype_original: {d['dtype_original']}")
print(f" slice_spec: {d['slice_spec']}")
print(f" slice shape={tuple(s.shape)} mean={s.mean():+.4f} std={s.std():.4f} min={s.min():+.4f} max={s.max():+.4f}")
print(f" latent shape={tuple(L.shape)} mean={L.mean():+.4f} std={L.std():.4f} min={L.min():+.4f} max={L.max():+.4f}")
print(f" finite: latent NaN={torch.isnan(L).any().item()} Inf={torch.isinf(L).any().item()}; "
f"slice NaN={torch.isnan(s).any().item()} Inf={torch.isinf(s).any().item()}")
print(f" metadata: {d['metadata']}\n")
```
Sanity criteria:
- `format_version == 1` (matches `LATENT_REFERENCE_FORMAT_VERSION`).
- `shape` matches what the model produces (e.g. LTX-2 distilled =
`[1, 128, T_lat, H_lat, W_lat]`; Stable Audio Open 1.0 = `[1, 64, 1024]`).
- `slice_spec.kind` matches a registered kind (`corner_3x3_first_frame`
for video, `audio_first_8_timesteps` for audio).
- No `NaN`/`Inf`. `mean ≈ 0`, `std ≈ 1` (denoised latents stay close to
the initial Gaussian distribution; very wide deviations suggest
numerical drift).
- `metadata.prompt` matches the test's prompt.
Then ask:
> "Numerics look right? Reply **`upload`** to continue, or anything else to abort."
Do not proceed until the user explicitly says `upload`. If they abort, leave
everything on disk so they can inspect further — no cleanup.
### 5. Copy into the local reference layout
Scoped copy — only the new test's artefacts. Single command works for both
artefact types because `_iter_reference_files` walks `.mp4` and `.pt`:
```bash
python fastvideo/tests/ssim/reference_videos_cli.py copy-local \
--quality-tier default \
--device-folder L40S_reference_videos \
--generated-dir ./generated_videos_modal/default/generated_videos/L40S_reference_videos
```
(The `--generated-dir` points at the device-folder root inside the
downloaded tree; `copy-local` walks all `<model>/<backend>/*.{mp4,pt}`
underneath it. Since the Modal run was scoped to a single test file via
`--test-files`, only that test's model(s) are present — so the copy is
implicitly per-test.)
Result for pixel: `fastvideo/tests/ssim/reference_videos/default/L40S_reference_videos/<model_id>/<backend>/<prompt>.mp4`.
Result for latent: same path with `.pt` extension.
### 6. Upload to HF — scoped per model_id, with overwrite guard
For each `<model_id>`:
```bash
python fastvideo/tests/ssim/reference_videos_cli.py upload \
--quality-tier default \
--device-folder L40S_reference_videos \
--model-id "<model_id>"
```
The upload command:
- Uploads **only** `reference_videos/default/L40S_reference_videos/<model_id>/`.
- **Refuses** if any file already exists at that path on HF (this is the
guard — seeding a new test should never clobber existing refs). To override,
the user must re-run with `--force`. If the guard fires, stop and report
exactly which files exist; do not silently `--force`.
Reads the HF token from `HF_API_KEY` / `HUGGINGFACE_HUB_TOKEN` / `HF_TOKEN`.
### 7. Report success
List what was uploaded (paths in repo) and remind the user to push any
related code changes. Do **not** auto-verify by re-running Modal — the user
can run `pytest fastvideo/tests/ssim/<test_file>` later to confirm end-to-end;
it will auto-download the refs they just uploaded.
## Failure modes and how to handle them
- **`HF_API_KEY` unset.** Stop before step 2. The Modal run needs it (passed
via `--hf-api-key`), and step 6 needs it for upload. If the user
ran `hf auth login` instead of exporting an env var, read the cached
token via `huggingface_hub.get_token()` and forward it to Modal as
`--hf-api-key="$CACHED_TOKEN"`.
- **Modal run fails before generation.** No artefacts on the volume — nothing
to download. Fix the test locally (`pytest fastvideo/tests/ssim/<test_file>`)
and retry from step 2.
- **`./generated_videos_modal/default/L40S_reference_videos/` missing after
`modal volume get`.** The run didn't produce artefacts (most likely the
test crashed before writing, or `REQUIRED_GPUS` exceeded the partition
capacity — see Modal logs).
- **Latent test crashed with FSDP / inference_mode error
(`RuntimeError: Inference tensors do not track version counter`).** The
test must pass `init_kwargs_override={"use_fsdp_inference": False}` when
`sp_size == 1` — see `test_stable_audio_similarity.py` for the pattern.
Fix in the test, push, retry.
- **Upload guard fires (files already exist).** The test name / model id
collides with something already on HF. Verify the user actually wants to
replace existing refs; if so, re-run the upload with `--force`. If not,
rename the model id in `*_MODEL_TO_PARAMS` and re-seed.
- **Quality looks wrong in step 4.** Abort. The artefacts stay on disk for
inspection. The fix is usually in the test's params (resolution, steps,
seed) — edit the test, then re-run the skill.
- For latent: also check `slice_spec.kind` matches the latent rank
(`corner_3x3_first_frame` requires 5-D, `audio_first_8_timesteps`
requires 3-D); a rank/kind mismatch raises in `_extract_expected_slice`.
## Design notes (for future skill maintainers)
- The skill deliberately runs on Modal, **not** locally, because the CI
runner is L40S. Seeding from a different GPU SKU produces refs that CI's
L40S runs can't match (pixel SSIM drifts across SKUs; latent cosine has
tighter cross-SKU bf16 drift but the configured tolerances assume
same-SKU seed → same-SKU verify).
- The skill is default-tier only. `full_quality` refs are seeded by a
separate, deliberate operation — they double runtime and aren't what CI
gates on.
- The overwrite guard in `reference_videos_cli.py upload` is default-on
specifically because this skill exists. Re-seeding is a distinct operation
that requires explicit `--force`.
- Both artefact types share the same Modal flow: the orchestrator sets
`--skip-reference-download` + `--no-fail-fast`, runs pytest, the test's
helper writes the artefact (`.mp4` via `imageio` for pixel,
`save_latent_reference` → `torch.save` for latent) BEFORE the
missing-reference assertion raises. `_sync_generated_videos_to_volume` in
`ssim_test.py` does a `shutil.copytree` of the whole `generated_videos/`
tree, picking up `.mp4`, `.pt`, and the `*_ssim.json` / `*_latent.json`
metric files alongside.
## References
- `fastvideo/tests/modal/ssim_test.py` — Modal orchestrator; see
`--sync-generated-to-volume`, `--generated-volume-subdir`,
`--skip-reference-download`, `--no-fail-fast`.
- `fastvideo/tests/ssim/reference_videos_cli.py` — `copy-local`, `upload`
(with `--model-id`, `--force`), `download`, `ensure` subcommands.
Extension allowlist is `REFERENCE_EXTENSIONS = VIDEO_EXTENSIONS +
LATENT_EXTENSIONS` (`.pt`).
- `fastvideo/tests/ssim/README.md` — reference layout, HF repo conventions.
- `fastvideo/tests/ssim/inference_similarity_utils.py` — pixel helpers
(`run_text_to_video_similarity_test`,
`run_image_to_video_similarity_test`, `build_init_kwargs`).
- `fastvideo/tests/ssim/latent_similarity_utils.py` — latent helper
(`run_text_to_latent_similarity_test`), slice spec dispatch
(`_extract_expected_slice`), reference schema
(`save_latent_reference` / `load_latent_reference`),
`LATENT_REFERENCE_FORMAT_VERSION`.
## Changelog
| Date | Change |
|------|--------|
| 2026-04-17 | Initial version (Modal sync-to-volume flow). |
| 2026-04-21 | Rewrite: single-test scope, explicit user-review pause, per-`model_id` upload, HF overwrite guard. Dropped `scripts/seed_ssim.sh`. |
| 2026-04-21 | Post-first-run fixes: `modal volume get` needs `--force` when parent exists; download tree has an extra `generated_videos/` level so `--generated-dir` must reflect it. |
| 2026-05-01 | Latent (`*.pt`) artefact support: artefact-type detection in step 1, type-aware verification (visual eyeball for mp4, numerics dump for pt) in step 4, FSDP+inference_mode failure-mode added, design notes for the unified Modal flow. Triggered by PR #1253 (LTX-2 latent migration + Stable Audio latent test). |
@@ -29,8 +29,8 @@
"Will Smith casually eats noodles, his relaxed demeanor contrasting with the energetic background of a bustling street food market. The scene captures a mix of humor and authenticity. Mid-shot framing, vibrant lighting."
],
"run_config": {
"num_warmup_runs": 1,
"num_measurement_runs": 3,
"num_warmup_runs": 2,
"num_measurement_runs": 5,
"required_gpus": 2
},
"thresholds": {
+88 -4
View File
@@ -15,8 +15,21 @@ log "Project root: $PROJECT_ROOT"
# Install Modal if not available
if ! python3 -m modal --version &> /dev/null; then
log "Modal not found, installing..."
python3 -m pip install modal
if ! command -v uv &> /dev/null; then
log "uv not found, bootstrapping..."
if ! curl -LsSf https://astral.sh/uv/install.sh | sh; then
log "Error: Failed to bootstrap uv via astral.sh installer."
exit 1
fi
export PATH="$HOME/.local/bin:$PATH"
if ! command -v uv &> /dev/null; then
log "Error: uv still not on PATH after bootstrap."
exit 1
fi
fi
# --break-system-packages preserves prior `pip install --user` semantics on PEP 668 agents.
uv pip install --system --break-system-packages modal
# Verify installation
if ! python3 -m modal --version &> /dev/null; then
log "Error: Failed to install modal. Please install it manually."
@@ -63,7 +76,72 @@ EFFECTIVE_PR=${BUILDKITE_PULL_REQUEST:-false}
if [ "$EFFECTIVE_PR" = "false" ] && [ -n "${PR_NUMBER:-}" ]; then
EFFECTIVE_PR=$PR_NUMBER
fi
MODAL_ENV="BUILDKITE_REPO=$BUILDKITE_REPO BUILDKITE_COMMIT=$BUILDKITE_COMMIT BUILDKITE_PULL_REQUEST=$EFFECTIVE_PR IMAGE_VERSION=$IMAGE_VERSION"
MODAL_ENV="BUILDKITE_REPO=$BUILDKITE_REPO BUILDKITE_COMMIT=$BUILDKITE_COMMIT BUILDKITE_PULL_REQUEST=$EFFECTIVE_PR BUILDKITE_BRANCH=${BUILDKITE_BRANCH:-} TEST_SCOPE=${TEST_SCOPE:-} IMAGE_VERSION=$IMAGE_VERSION"
POST_RUN_HOOK=""
upload_performance_artifacts() {
SHORT_SHA=${BUILDKITE_COMMIT:0:7}
LOCAL_DIR="downloaded_reports"
_download_reports() {
log "Downloading perf_reports/ from Modal Volume..."
mkdir -p "$LOCAL_DIR"
if ! modal volume get hf-model-weights "perf_reports/" "$LOCAL_DIR"; then
log "Error: Failed to download perf_reports/ from Modal Volume."
return 1
fi
}
_upload_dashboard() {
local target
target=$(find "$LOCAL_DIR" -name "dashboard_${SHORT_SHA}_*" | head -n 1)
log "TARGET dashboard: '$target'"
if [ -n "$target" ]; then
log "Found dashboard: $target. Uploading to Buildkite..."
buildkite-agent artifact upload "$target"
buildkite-agent annotate --style info --context "perf-dashboard" < "$target"
else
log "Warning: Could not find a dashboard file matching $SHORT_SHA"
fi
}
_upload_perf_summary() {
local target
target=$(find "$LOCAL_DIR" -name "perf_${SHORT_SHA}_*" | head -n 1)
log "TARGET perf summary: '$target'"
if [ -n "$target" ]; then
log "Found perf summary: $target. Uploading to Buildkite..."
buildkite-agent artifact upload "$target"
buildkite-agent annotate --style info --context "perf-summary" < "$target"
else
log "Warning: Could not find a perf summary file matching $SHORT_SHA"
fi
}
_cleanup_modal_volume() {
log "Cleaning up perf_reports/ from Modal Volume..."
if modal volume rm hf-model-weights "perf_reports/" --recursive; then
log "Successfully deleted perf_reports/ from Modal Volume."
else
log "Warning: Failed to delete perf_reports/ from Modal Volume. Manual cleanup may be required."
fi
}
_cleanup_local() {
log "Cleaning up local download directory..."
rm -rf "$LOCAL_DIR"
}
# --- Main flow ---
_download_reports || { _cleanup_local; return 1; }
_upload_dashboard
_upload_perf_summary
_cleanup_modal_volume
_cleanup_local
}
case "$TEST_TYPE" in
"encoder")
@@ -124,8 +202,9 @@ case "$TEST_TYPE" in
MODAL_COMMAND="$MODAL_ENV HF_API_KEY=$HF_API_KEY python3 -m modal run $MODAL_TEST_FILE::run_lora_extraction_tests"
;;
"performance")
log "Running performance tests..."
log "Running performance tests on Modal..."
MODAL_COMMAND="$MODAL_ENV HF_API_KEY=$HF_API_KEY python3 -m modal run $MODAL_TEST_FILE::run_performance_tests"
POST_RUN_HOOK="upload_performance_artifacts"
;;
"api_server")
log "Running API server integration tests..."
@@ -147,5 +226,10 @@ else
log "Error: Modal test failed with exit code: $TEST_EXIT_CODE"
fi
if [ -n "$POST_RUN_HOOK" ]; then
log "Executing post-run hook: $POST_RUN_HOOK"
"$POST_RUN_HOOK"
fi
log "=== Test execution completed with exit code: $TEST_EXIT_CODE ==="
exit $TEST_EXIT_CODE
+15 -2
View File
@@ -13,8 +13,21 @@ log "Project root: $PROJECT_ROOT"
if ! python3 -m pre_commit --version &> /dev/null; then
log "pre-commit not found, installing..."
python3 -m pip install --user pre-commit==4.0.1
if ! command -v uv &> /dev/null; then
log "uv not found, bootstrapping..."
if ! curl -LsSf https://astral.sh/uv/install.sh | sh; then
log "Error: Failed to bootstrap uv via astral.sh installer."
exit 1
fi
export PATH="$HOME/.local/bin:$PATH"
if ! command -v uv &> /dev/null; then
log "Error: uv still not on PATH after bootstrap."
exit 1
fi
fi
# --break-system-packages preserves prior `pip install --user` semantics on PEP 668 agents.
uv pip install --system --break-system-packages pre-commit==4.0.1
if ! python3 -m pre_commit --version &> /dev/null; then
log "Error: Failed to install pre-commit."
exit 1
+4 -3
View File
@@ -37,10 +37,11 @@ jobs:
with:
python-version: '3.12'
- name: Install uv
uses: astral-sh/setup-uv@v3
- name: Install dependencies
run: |
python -m pip install --upgrade pip
pip install -r requirements-mkdocs.txt
run: uv pip install --system -r requirements-mkdocs.txt
- name: Setup Pages
uses: actions/configure-pages@v4
+4 -3
View File
@@ -56,10 +56,11 @@ jobs:
with:
python-version: '3.10'
- name: Install uv
uses: astral-sh/setup-uv@v3
- name: Install build dependencies
run: |
python -m pip install --upgrade pip
pip install build twine wheel
run: uv pip install --system build twine wheel
- name: Build package
run: |
+16 -11
View File
@@ -131,11 +131,13 @@ jobs:
clang-11 --version
nvcc --version
- name: Install uv
uses: astral-sh/setup-uv@v3
- name: Install PyTorch ${{ matrix.torch-cuda.torch-version }}+cu${{ matrix.torch-cuda.cuda-version }}
run: |
pip install --upgrade pip
pip install typing-extensions==4.12.2
pip install --no-cache-dir torch==${{ matrix.torch-cuda.torch-version }} --index-url https://download.pytorch.org/whl/${{matrix.torch-cuda.torch-cuda-short}}
uv pip install --system typing-extensions==4.12.2
uv pip install --system --no-cache-dir torch==${{ matrix.torch-cuda.torch-version }} --index-url https://download.pytorch.org/whl/${{matrix.torch-cuda.torch-cuda-short}}
nvcc --version
python --version
python -c "import torch; print('PyTorch:', torch.__version__)"
@@ -145,20 +147,20 @@ jobs:
- name: Build wheel
run: |
export PYTHONPATH=$GITHUB_WORKSPACE:$PYTHONPATH
pip install setuptools ninja packaging wheel triton scikit-build-core cmake build
uv pip install --system setuptools ninja packaging wheel triton scikit-build-core cmake build
cd fastvideo-kernel
git submodule update --init --recursive # Ensure ThunderKittens submodule is initialized
# Release builds are produced on GPU-less runners, so force-enable TK and target Hopper.
export TORCH_CUDA_ARCH_LIST="9.0a"
export CMAKE_ARGS="${CMAKE_ARGS:-} -DFASTVIDEO_KERNEL_BUILD_TK=ON -DCMAKE_CUDA_ARCHITECTURES=90a"
# Build standard wheel (no local version suffix) for PyPI
python -m build --wheel --outdir dist
# Fix the wheel to be manylinux compliant
pip install auditwheel
uv pip install --system auditwheel
# Point auditwheel at torch libs, but do not vendor them into the wheel.
TORCH_LIB_DIR=$(python - <<'PY'
import os
@@ -211,10 +213,13 @@ jobs:
pattern: 'fastvideo_kernel-py*'
merge-multiple: true
- name: Install uv
uses: astral-sh/setup-uv@v3
- name: Build source distribution
run: |
pip install build scikit-build-core cmake ninja
uv pip install --system build scikit-build-core cmake ninja
cd fastvideo-kernel
# We don't need full CUDA/Torch to just package the source (sdist)
python -m build --sdist --outdir dist
+8
View File
@@ -92,3 +92,11 @@ preprocess_output_text/
.sisyphus/
openspec/
fastvideo/tests/ssim/reference_videos/**
# Local clones of upstream repos used only for parity testing.
/stable-audio-tools/
/daVinci-MagiHuman/
# Converted model weights (produced by scripts/checkpoint_conversion/*).
# Tens of GB; should live on HF, not in git.
/converted_weights/
+2 -9
View File
@@ -7,20 +7,13 @@ exclude: |
fastvideo-kernel/.*|
assets/.*|
tests/.*|
demo/.*|
predict\.py|
scripts/.*|
assets/prompts/.*|
fastvideo/data_preprocess/.*|
fastvideo/dataset/.*|
fastvideo/models/.*|
fastvideo/sample/.*|
fastvideo/train\.py|
fastvideo/utils/.*|
examples/.*|
\.agents/.*|
.github/workflows/publish-fastvideo.yml|
.github/workflows/_template-build-image.yml|
docs/source/inference/support_matrix.md
.github/workflows/_template-build-image.yml
)
repos:
- repo: https://github.com/google/yapf
+31 -2
View File
@@ -11,7 +11,7 @@
- Static assets: `assets/` (including `assets/images/`, `assets/videos/`, and `assets/prompts/`) and `comfyui/assets/`.
## Build, Test, and Development Commands
- `uv pip install -e .[dev]`: editable install with lint/test extras.
- `uv pip install -e ".[dev]"`: editable install with lint/test extras.
- `pre-commit install --hook-type pre-commit --hook-type commit-msg`: enable local hooks.
- `pre-commit run --all-files`: run formatter/lint/type/spelling checks.
- `pytest tests/`: run top-level test suite.
@@ -23,7 +23,8 @@
- Python 3.10+; 4-space indentation; keep code and imports readable and explicit.
- Style tools are configured in `pyproject.toml` and `.pre-commit-config.yaml`:
- `yapf` (format), `ruff` (lint, auto-fix), `mypy` (typing), `codespell`.
- Target line length is 80.
- Lint via `pre-commit run --files <changed paths>` (or `pre-commit run --all-files` for a full sweep) before committing. Do not shell out to `yapf`/`ruff`/`codespell`/`mypy` directly — pre-commit chains them with the project's config and respects the `.pre-commit-config.yaml` excludes (e.g. `fastvideo/tests/` is intentionally skipped). If pre-commit reports `(no files to check)` for your paths, that exclude is deliberate — don't bypass it.
- Target line length is 120 (configured in `pyproject.toml` for ruff, yapf, and isort).
- Naming: `snake_case` for functions/files, `PascalCase` for classes, `UPPER_SNAKE_CASE` for constants.
## Testing Guidelines
@@ -54,3 +55,31 @@ This repository is agent-friendly. Before doing any work, read:
If you are exploring a new procedure that has no existing SOP, document your
progress in `.agents/exploration/` and flag it for review at the end of your
session.
## Per-Directory AGENTS.md
Local guidance lives next to the code. Read the in-scope file before editing:
| Directory | What it covers |
|-----------|----------------|
| `fastvideo/AGENTS.md` | Core package map, public API, registry-driven model dispatch |
| `fastvideo/configs/AGENTS.md` | Arch + pipeline config dataclasses, `param_names_mapping` |
| `fastvideo/models/AGENTS.md` | DiT / VAE / encoder / scheduler / loader layout (pre-commit excluded) |
| `fastvideo/layers/AGENTS.md` | Tensor-parallel linear/attention layer rules for ports |
| `fastvideo/attention/AGENTS.md` | Backend registry + env-var override |
| `fastvideo/pipelines/AGENTS.md` | Stage ABC, `basic/<model>/`, `preprocess/`, presets |
| `fastvideo/training/AGENTS.md` | Legacy monolithic pipelines (frozen for existing models) |
| `fastvideo/train/AGENTS.md` | New modular trainer (methods × models × callbacks, YAML) |
| `fastvideo/tests/AGENTS.md` | Test taxonomy, conftest, pre-commit-excluded path |
| `fastvideo/tests/ssim/AGENTS.md` | GPU SSIM regression authoring + reference video sync |
| `scripts/checkpoint_conversion/AGENTS.md` | Adding a converter for a new HF/official checkpoint |
## Critical: Two Training Stacks Coexist
- `fastvideo/training/` — legacy, monolithic per-model `*_training_pipeline.py` and
`*_distillation_pipeline.py`. Still authoritative for shipped models.
- `fastvideo/train/` — new modular framework (composable methods × models × callbacks
driven by YAML). Preferred for new training work.
Pick the matching stack before editing. Do not migrate a pipeline between them
without an explicit ask — the conventions and config surfaces differ.
+2 -2
View File
@@ -128,7 +128,7 @@ class CLIPFeatureExtractor(BaseFeatureExtractor):
def __init__(self, device: str = 'cuda', model_name: str = "openai/clip-vit-base-patch32"):
if not TRANSFORMERS_AVAILABLE:
raise ImportError("Please install transformers: pip install transformers")
raise ImportError("Please install transformers: uv pip install transformers")
super().__init__(device)
self.processor = CLIPProcessor.from_pretrained(model_name)
self.model = CLIPModel.from_pretrained(model_name).to(self.device)
@@ -171,7 +171,7 @@ class VideoMAEFeatureExtractor(BaseFeatureExtractor):
def __init__(self, device: str = 'cuda', model_name: str = "MCG-NJU/videomae-base"):
if not TRANSFORMERS_AVAILABLE:
raise ImportError("Please install transformers: pip install transformers")
raise ImportError("Please install transformers: uv pip install transformers")
super().__init__(device)
self.model = VideoMAEModel.from_pretrained(model_name).to(self.device)
self.model.eval()
+1 -1
View File
@@ -57,7 +57,7 @@ class I3DFeatureExtractor(nn.Module):
except Exception as e:
raise RuntimeError(f"Failed to load I3D model from Hugging Face Hub. Error: {e}\n"
f"Ensure you have internet connection and huggingface_hub installed:\n"
f"pip install huggingface_hub") from e
f"uv pip install huggingface_hub") from e
def preprocess(self, videos: torch.Tensor) -> torch.Tensor:
"""
+1 -1
View File
@@ -1,7 +1,7 @@
#!/bin/bash
# 1. Install missing dependency
pip install -q opencv-python-headless transformers huggingface_hub
uv pip install -q opencv-python-headless transformers huggingface_hub
# 2. Run FVD script
python benchmarks/fvd/run_fvd.py
+1 -1
View File
@@ -1,4 +1,4 @@
#!/bin/bash
# 1. Install missing dependency
pip install -q opencv-python-headless
uv pip install -q opencv-python-headless
+2 -2
View File
@@ -38,10 +38,10 @@ cp -r /path/to/FastVideo/comfyui /path/to/ComfyUI/custom_nodes/FastVideo
#### Install dependencies:
Currently, the only dependency is `fastvideo`, which can be installed using pip.
Currently, the only dependency is `fastvideo`, which can be installed with `uv`.
```bash
pip install fastvideo
uv pip install fastvideo
```
#### Install missing custom nodes:
+3 -3
View File
@@ -42,15 +42,15 @@ RUN source $HOME/.local/bin/env && \
uv venv --python 3.10 --seed /opt/venv && \
source /opt/venv/bin/activate && \
uv pip install --no-cache-dir --upgrade pip && \
uv pip install --no-cache-dir .[dev] && \
uv pip install --no-cache-dir https://github.com/mjun0812/flash-attention-prebuild-wheels/releases/download/v0.7.16/flash_attn-2.8.3+cu128torch2.10-cp310-cp310-linux_x86_64.whl
uv pip install --no-cache-dir ".[dev]" && \
uv pip install --no-cache-dir https://github.com/mjun0812/flash-attention-prebuild-wheels/releases/download/v0.9.4/flash_attn-2.8.3+cu128torch2.11-cp310-cp310-linux_x86_64.whl
COPY . .
# Install dependencies using uv and set up shell configuration
RUN source $HOME/.local/bin/env && \
source /opt/venv/bin/activate && \
uv pip install --no-cache-dir -e .[dev] && \
uv pip install --no-cache-dir -e ".[dev]" && \
git config --unset-all http.https://github.com/.extraheader || true && \
echo 'source /opt/venv/bin/activate' >> /root/.bashrc && \
echo 'if [ -n "$ZSH_VERSION" ] && [ -f ~/.zshrc ]; then . ~/.zshrc; elif [ -f ~/.bashrc ]; then . ~/.bashrc; fi' > /root/.profile
+3 -3
View File
@@ -42,15 +42,15 @@ RUN source $HOME/.local/bin/env && \
uv venv --python 3.11 --seed /opt/venv && \
source /opt/venv/bin/activate && \
uv pip install --no-cache-dir --upgrade pip && \
uv pip install --no-cache-dir .[dev] && \
uv pip install --no-cache-dir https://github.com/mjun0812/flash-attention-prebuild-wheels/releases/download/v0.7.16/flash_attn-2.8.3+cu128torch2.10-cp311-cp311-linux_x86_64.whl
uv pip install --no-cache-dir ".[dev]" && \
uv pip install --no-cache-dir https://github.com/mjun0812/flash-attention-prebuild-wheels/releases/download/v0.9.4/flash_attn-2.8.3+cu128torch2.11-cp311-cp311-linux_x86_64.whl
COPY . .
# Install dependencies using uv and set up shell configuration
RUN source $HOME/.local/bin/env && \
source /opt/venv/bin/activate && \
uv pip install --no-cache-dir -e .[dev] && \
uv pip install --no-cache-dir -e ".[dev]" && \
git config --unset-all http.https://github.com/.extraheader || true && \
echo 'source /opt/venv/bin/activate' >> /root/.bashrc && \
echo 'if [ -n "$ZSH_VERSION" ] && [ -f ~/.zshrc ]; then . ~/.zshrc; elif [ -f ~/.bashrc ]; then . ~/.bashrc; fi' > /root/.profile
+3 -3
View File
@@ -42,15 +42,15 @@ RUN source $HOME/.local/bin/env && \
uv venv --python 3.12 --seed /opt/venv && \
source /opt/venv/bin/activate && \
uv pip install --no-cache-dir --upgrade pip && \
uv pip install --no-cache-dir .[dev] && \
uv pip install --no-cache-dir https://github.com/mjun0812/flash-attention-prebuild-wheels/releases/download/v0.7.16/flash_attn-2.8.3+cu128torch2.10-cp312-cp312-linux_x86_64.whl
uv pip install --no-cache-dir ".[dev]" && \
uv pip install --no-cache-dir https://github.com/mjun0812/flash-attention-prebuild-wheels/releases/download/v0.9.4/flash_attn-2.8.3+cu128torch2.11-cp312-cp312-linux_x86_64.whl
COPY . .
# Install dependencies using uv and set up shell configuration
RUN source $HOME/.local/bin/env && \
source /opt/venv/bin/activate && \
uv pip install --no-cache-dir -e .[dev] && \
uv pip install --no-cache-dir -e ".[dev]" && \
git config --unset-all http.https://github.com/.extraheader || true && \
echo 'source /opt/venv/bin/activate' >> /root/.bashrc && \
echo 'if [ -n "$ZSH_VERSION" ] && [ -f ~/.zshrc ]; then . ~/.zshrc; elif [ -f ~/.bashrc ]; then . ~/.bashrc; fi' > /root/.profile
+2 -2
View File
@@ -42,7 +42,7 @@ RUN source $HOME/.local/bin/env && \
uv venv --python 3.12 --seed /opt/venv && \
source /opt/venv/bin/activate && \
uv pip install --no-cache-dir --upgrade pip && \
uv pip install --no-cache-dir .[dev] && \
uv pip install --no-cache-dir ".[dev]" && \
uv pip install --no-cache-dir flash-attn==2.8.3 --no-build-isolation
COPY . .
@@ -50,7 +50,7 @@ COPY . .
# Install dependencies using uv and set up shell configuration
RUN source $HOME/.local/bin/env && \
source /opt/venv/bin/activate && \
uv pip install --no-cache-dir -e .[dev] && \
uv pip install --no-cache-dir -e ".[dev]" && \
git config --unset-all http.https://github.com/.extraheader || true && \
echo 'source /opt/venv/bin/activate' >> /root/.bashrc && \
echo 'if [ -n "$ZSH_VERSION" ] && [ -f ~/.zshrc ]; then . ~/.zshrc; elif [ -f ~/.bashrc ]; then . ~/.bashrc; fi' > /root/.profile
+1 -1
View File
@@ -43,7 +43,7 @@ COPY . .
# Install dependencies using uv and set up shell configuration
RUN source $HOME/.local/bin/env && \
source /opt/venv/bin/activate && \
uv pip install --no-cache-dir -e .[rocm] && \
uv pip install --no-cache-dir -e ".[rocm]" && \
git config --unset-all http.https://github.com/.extraheader || true && \
echo 'source /opt/venv/bin/activate' >> /root/.bashrc && \
echo 'if [ -n "$ZSH_VERSION" ] && [ -f ~/.zshrc ]; then . ~/.zshrc; elif [ -f ~/.bashrc ]; then . ~/.bashrc; fi' > /root/.profile
+1 -1
View File
@@ -6,7 +6,7 @@ This directory contains the FastVideo documentation built with MkDocs.
```bash
# Install dependencies
pip install -r requirements-mkdocs.txt
uv pip install -r requirements-mkdocs.txt
# Serve docs with live reload (recommended for development)
mkdocs serve
+210
View File
@@ -0,0 +1,210 @@
# Activation Trace Mode
!!! note
This page covers Extension 0 (module forward hooks), which is the implemented
tracing mechanism. Extensions 1-3 are design sketches for future work and are
**not yet implemented**.
## Overview
Activation trace mode is a zero-overhead-when-off, env-gated mechanism for
dumping per-layer activation statistics during FastVideo inference. Its primary
use case is **parity debugging across model ports**: enable tracing on both
FastVideo and the upstream reference implementation, then `diff` the resulting
JSONL files to find the first divergent layer.
The mechanism is intentionally narrow. It doesn't replace general logging,
profiling, or function tracing. It answers one question: "at which layer do
FastVideo and the reference model first produce different numbers?"
## When to use
- Investigating numerical drift between FastVideo and an upstream reference.
- Debugging mid-pipeline divergence (e.g., one block produces wrong output while earlier blocks match).
- Validating that a refactor preserves bf16 noise-floor behavior across many layers.
## When NOT to use
| Goal | Use instead |
|---|---|
| General logging | `init_logger(__name__)` |
| Per-stage timing | `FASTVIDEO_STAGE_LOGGING` |
| Profiling kernel timings | `FASTVIDEO_TORCH_PROFILER_DIR` (see [Profiling](profiling.md)) |
| Function-call tracing | `FASTVIDEO_TRACE_FUNCTION` (heavy) |
## Quickstart
```bash
FASTVIDEO_TRACE_ACTIVATIONS=1 \
FASTVIDEO_TRACE_LAYERS="^block\.layers\.[0-9]+$" \
FASTVIDEO_TRACE_STATS="abs_mean,sum,max,shape" \
FASTVIDEO_TRACE_OUTPUT="/tmp/fv_trace.jsonl" \
python examples/inference/basic/basic_magi_human.py
```
Each line in `/tmp/fv_trace.jsonl` is a JSON record:
```json
{"module": "block.layers.0", "tensor": "out", "step": 0, "abs_mean": 1.234, "sum": -5.678, "max": 9.012, "shape": [1, 4096, 5120]}
```
## Configuration
| Env var | Default | Description |
|---|---|---|
| `FASTVIDEO_TRACE_ACTIVATIONS` | `False` | Master toggle. When unset or false, **zero overhead** in the production hot path. |
| `FASTVIDEO_TRACE_LAYERS` | `""` (all) | Python regex filter applied to `model.named_modules()` names. Empty string matches all modules. |
| `FASTVIDEO_TRACE_STATS` | `"abs_mean,sum"` | Comma-separated stats to compute. Available: `abs_mean`, `sum`, `min`, `max`, `mean`, `std`, `shape`, `dtype`. |
| `FASTVIDEO_TRACE_OUTPUT` | `"/tmp/fv_trace_<pid>.jsonl"` | Output file path. `<pid>` is replaced with the process ID at runtime. |
| `FASTVIDEO_TRACE_STEPS` | `""` (all) | Comma-separated denoising step indices to capture. Empty string captures all steps. |
## Workflow: parity-debug a model port
1. Set up a tightly-controlled comparison: a parity test or a small standalone
script that loads both the FastVideo model and the upstream reference with
identical inputs and seeds.
2. Run the FastVideo side with tracing on:
```bash
FASTVIDEO_TRACE_ACTIVATIONS=1 \
FASTVIDEO_TRACE_LAYERS="<your regex>" \
FASTVIDEO_TRACE_OUTPUT="/tmp/fv_trace_fv.jsonl" \
python <fv_runner.py>
```
3. Run the upstream side. The upstream repo needs separate instrumentation. See
"Hooking the upstream side" below.
4. Sort both files by `(module, step)` if needed, then diff:
```bash
diff /tmp/fv_trace_fv.jsonl /tmp/fv_trace_upstream.jsonl
```
5. The first divergent line identifies the first layer where FastVideo and the
upstream produce different outputs. Start debugging there.
## Architecture (Extension 0: module forward hooks)
At pipeline initialization, `attach_activation_trace()` reads the env vars once.
If `FASTVIDEO_TRACE_ACTIVATIONS` is unset or false, the function returns
immediately and no hooks are registered. If tracing is on, it walks
`model.named_modules()`, filters by the layer regex, and registers an
`ActivationStatHook` on each matching module.
During the forward pass, each hook fires after its module completes, computes
the requested stats on the output tensor, and appends a JSON record to the
output file.
```
ComposedPipelineBase
└─ attach_activation_trace()
├─ reads env vars (once at startup)
├─ if off: returns None immediately
└─ if on: walks named_modules()
└─ registers ActivationStatHook on matching modules
└─ on each forward: compute stats → append JSONL
```
### Zero-overhead-when-off guarantee
- The env var check happens **once at startup** inside `attach_activation_trace()`.
- If the env var is unset or false, the function returns `None` immediately.
- No hooks are registered. No branches are added to the production forward path.
- The only cost when tracing is off is one env var lookup at pipeline
initialization, which takes under a microsecond.
### Hooking the upstream side
The upstream reference repo isn't part of FastVideo, so it can't read FastVideo
env vars directly. Two options:
**Option 1: Inline patch** in your local clone of the upstream repo. Add
`register_forward_hook` calls in the same shape as `ActivationStatHook`. Clean
up afterward with `git stash` or `git checkout HEAD -- <file>`.
**Option 2: Wrapper script**. Write a small Python harness that imports the
upstream model, walks its `named_modules()`, and attaches hooks externally.
This is the same pattern used in
`tests/local_tests/transformers/_debug_magi_human_block_parity.py`.
The `add-model-trace` skill at `~/.config/opencode/skill/add-model-trace/`
provides a script template for this purpose.
## Future extensions (design only, not yet implemented)
### Extension 1: FX/Dynamo backend graph rewrite
**Granularity**: per-FX-node (every matmul, every add).
**Mechanism**: a `torch.compile` backend that takes the captured `GraphModule`
and inserts logger nodes after each op. Compiles into a separate artifact from
the production graph.
**Off semantics**: zero overhead. The production compile path is untouched.
**When to add**: if you need to trace inside a `torch.compile`'d graph and
Extension 0 is too coarse.
**Build cost**: roughly 1-2 days. Reference:
`torchao.quantization.pt2e._numeric_debugger`.
### Extension 2: AST source injection at import time
**Granularity**: per-line (between any two Python statements).
**Mechanism**: an importlib loader hook rewrites Python source AST at module
import time, inserting `if TRACE: dump(...)` statements. The decision is made
once at import.
**Off semantics**: zero overhead. If the env var is off at import time, source
is loaded as-is.
**When to add**: if you need per-line granularity that even FX-node-level can't
provide. This is almost never the right choice.
**Build cost**: roughly 1 week. Brittle and hard to debug.
### Extension 3: `__torch_dispatch__` / `TorchDispatchMode`
**Granularity**: per-op (every dispatcher call: matmul, add, view, etc.).
**Mechanism**: a `TorchDispatchMode` context manager that intercepts all ops at
the dispatcher level.
**Off semantics**: zero overhead. PyTorch's dispatcher only invokes mode hooks
when a mode is active.
**When on**: significant overhead. Every op pays a Python callback cost. Triton
kernels bypass it.
**When to add**: useful for quantization or dtype debugging where module-level
granularity isn't enough.
**Build cost**: roughly 1 day. Reference:
`torch.utils._python_dispatch.TorchDispatchMode`.
## Comparison with similar tools
| Tool | Pattern | FastVideo equivalent |
|---|---|---|
| SGLang `--debug-tensor-dump-output-folder` | env-gated forward hooks at startup | Extension 0 (this) |
| TransformerEngine `DumpTensors` | config-driven selective dumps | Extension 0 (env-driven) |
| HuggingFace `output_hidden_states=True` | source-level boolean gating | Not used; Extension 0 avoids model code edits |
| torchao numeric debugger | FX pass + node-level loggers | Extension 1 (future) |
| W&B `wandb.watch()` | runtime forward hooks (always on once registered) | Extension 0 has a similar mechanism, but gated off by default |
## Implementation references
- Module: `fastvideo/hooks/activation_trace.py`
- Env vars: `fastvideo/envs.py` (`FASTVIDEO_TRACE_ACTIVATIONS` and friends)
- Pipeline integration: `fastvideo/pipelines/composed_pipeline_base.py`
- Tests: `fastvideo/tests/hooks/test_activation_trace.py`
- Companion skill (for ad-hoc port investigations): `~/.config/opencode/skill/add-model-trace/`
## Changelog
| Date | Change |
|---|---|
| 2026-05-01 | Initial Extension 0 (module forward hooks) implementation. Extensions 1-3 designed but not implemented. |
+6 -3
View File
@@ -296,8 +296,10 @@ Action:
- Add or reuse a numerical parity test that loads the official model and the
FastVideo model and compares outputs.
- See examples in `tests/local_tests/` (e.g., `tests/local_tests/upsamplers/`)
and the commands in `tests/local_tests/README.md`.
- See examples in `tests/local_tests/` organized by model family
(e.g., `tests/local_tests/sd35/`, `tests/local_tests/ltx2/`,
`tests/local_tests/stable_audio/`) and the navigation index in
`tests/local_tests/README.md`.
- If there are discrepancies, add opt‑in logging to both models and compare
activation summaries (layer output sums, per‑stage logs).
- First align the loaded weights (validate `param_names_mapping`).
@@ -348,7 +350,8 @@ Purpose:
Action:
- Add a pipeline parity test under `tests/local_tests/pipelines/`.
- Add a pipeline parity test under `tests/local_tests/<family>/`
(e.g., `tests/local_tests/<family>/test_<family>_pipeline_parity.py`).
- See the [Testing Guide](testing.md) for test conventions.
### 7) Add user‑facing examples
+1 -1
View File
@@ -99,7 +99,7 @@ cd /FastVideo
**Install the package**
```bash
uv pip install -e .[dev]
uv pip install -e ".[dev]"
```
The Docker image already includes Flash Attention and most heavy dependencies, so this is fast.
+1 -1
View File
@@ -49,7 +49,7 @@ git clone https://github.com/hao-ai-lab/FastVideo.git && cd FastVideo
Install FastVideo in editable mode and set up hooks:
```bash
uv pip install -e .[dev]
uv pip install -e ".[dev]"
# Optional: FlashAttention (builds native kernels)
uv pip install flash-attn --no-build-isolation -v
@@ -29,7 +29,7 @@ surfaces:
vae_cpu_offload: generator.engine.offload.vae
pin_cpu_memory: generator.engine.offload.pin_cpu_memory
enable_torch_compile: generator.engine.compile.enabled
torch_compile_kwargs: generator.engine.compile.kwargs
torch_compile_kwargs: generator.engine.compile.backend,fullgraph,mode,dynamic,extras
disable_autocast: generator.engine.disable_autocast
enable_stage_verification: generator.engine.enable_stage_verification
prompt_txt: request.inputs.prompt_path
@@ -40,8 +40,8 @@ surfaces:
init_weights_from_safetensors_2: generator.pipeline.components.transformer_2_weights
override_pipeline_cls_name: generator.pipeline.components.override_pipeline_cls_name
boundary_ratio: request.sampling.boundary_ratio
ltx2_vae_tiling: generator.pipeline.vae_tiling
preset_owned:
ltx2_vae_tiling: generator.pipeline.preset_overrides.ltx2.vae_tiling
ltx2_vae_spatial_tile_size_in_pixels: generator.pipeline.preset_overrides.ltx2.vae.spatial_tile_size_in_pixels
ltx2_vae_spatial_tile_overlap_in_pixels: generator.pipeline.preset_overrides.ltx2.vae.spatial_tile_overlap_in_pixels
ltx2_vae_temporal_tile_size_in_frames: generator.pipeline.preset_overrides.ltx2.vae.temporal_tile_size_in_frames
@@ -306,6 +306,30 @@ surfaces:
sources: [fastvideo.configs.pipelines.wan.MatrixGameI2V480PConfig]
num_frames_per_block:
sources: [fastvideo.configs.pipelines.wan.MatrixGameI2V480PConfig]
audio_channels:
sources:
- fastvideo.configs.pipelines.stable_audio.StableAudioT2AConfig
- fastvideo.configs.pipelines.stable_audio.StableAudioOpenSmallConfig
audio_end_in_s:
sources:
- fastvideo.configs.pipelines.stable_audio.StableAudioT2AConfig
- fastvideo.configs.pipelines.stable_audio.StableAudioOpenSmallConfig
audio_start_in_s:
sources:
- fastvideo.configs.pipelines.stable_audio.StableAudioT2AConfig
- fastvideo.configs.pipelines.stable_audio.StableAudioOpenSmallConfig
max_audio_duration_s:
sources:
- fastvideo.configs.pipelines.stable_audio.StableAudioT2AConfig
- fastvideo.configs.pipelines.stable_audio.StableAudioOpenSmallConfig
sample_size:
sources:
- fastvideo.configs.pipelines.stable_audio.StableAudioT2AConfig
- fastvideo.configs.pipelines.stable_audio.StableAudioOpenSmallConfig
sampling_rate:
sources:
- fastvideo.configs.pipelines.stable_audio.StableAudioT2AConfig
- fastvideo.configs.pipelines.stable_audio.StableAudioOpenSmallConfig
compatibility_only:
batch_size: "Gen3C inference-only tuning field pending typed batching design."
gradient_checkpointing: "Gen3C inference-only compatibility field pending typed batching design."
@@ -354,6 +378,8 @@ surfaces:
return_frames: request.output.return_frames
return_trajectory_latents: request.runtime.return_trajectory_latents
return_trajectory_decoded: request.runtime.return_trajectory_decoded
continuation_state: request.state
return_continuation_state: request.output.return_state
preset_owned:
t_thresh: request.stage_overrides.refine.t_thresh
spatial_refine_only: request.stage_overrides.refine.spatial_refine_only
@@ -378,6 +404,13 @@ surfaces:
ltx2_stg_scale_audio: request.extensions.ltx2.stg_scale_audio
ltx2_stg_blocks_video: request.extensions.ltx2.stg_blocks_video
ltx2_stg_blocks_audio: request.extensions.ltx2.stg_blocks_audio
audio_start_in_s: request.extensions.stable_audio.audio_start_in_s
audio_end_in_s: request.extensions.stable_audio.audio_end_in_s
init_audio: request.extensions.stable_audio.init_audio
init_audio_strength: request.extensions.stable_audio.init_audio_strength
init_noise_level: request.extensions.stable_audio.init_noise_level
inpaint_audio: request.extensions.stable_audio.inpaint_audio
inpaint_mask: request.extensions.stable_audio.inpaint_mask
internal_only:
data_type: "Derived from the request shape and not a public input."
+177
View File
@@ -0,0 +1,177 @@
# Streaming WebSocket Server Contract
The streaming server (`fastvideo/entrypoints/streaming/server.py`) speaks
a JSON-over-WebSocket protocol with binary fMP4 chunks for media. This
document is the authoritative spec for the message catalogue and the
session state machine. Any change to either must update this document
in the same PR that touches `protocol.py` or `session.py`.
## Endpoint
| Path | Protocol | Purpose |
|---|---|---|
| `WS /v1/stream` | WebSocket (JSON + binary) | Per-session realtime streaming |
| `GET /health` | HTTP | Liveness probe (`status`, `stream_mode`, active `sessions`) |
The server is launched by `fastvideo serve --config <serve.yaml>` when
the config carries a `streaming:` block. Without that block the same CLI
launches the OpenAI stateless HTTP server instead.
## Connection lifecycle
Every WebSocket connection holds exactly one `Session`. Sessions move
through the states in `SessionState` (`fastvideo/entrypoints/streaming/session.py`).
```
┌──────────────┐
│ INITIALIZING │ ← WebSocket accepted, before init frame
└──────┬───────┘
│ session_init_v2 received
┌──────────────┼──────────────┐
▼ ▼ ▼
QUEUED GPU_BINDING REJECTED
│ │ ↑
│ slot ready │ │ max-sessions hit
▼ ▼ │ or invalid init
┌────────┐ │
│ ACTIVE │ ────────┘
└────┬───┘
segment loop │
│
┌───────────┼───────────┐
▼ ▼ ▼
COMPLETE ERROR TIMEOUT
(clean leave) (any failure) (idle / segment_cap reached)
```
Terminal states (`COMPLETE`, `ERROR`, `TIMEOUT`, `REJECTED`) are sinks —
no transitions out. The transition matrix is enforced in
`session.py::_VALID_TRANSITIONS`; bad transitions raise.
`SessionManager` enforces the per-process budgets pulled from
`StreamingConfig`:
- `session_timeout_seconds` — idle reaper drops sessions that haven't
advanced; non-terminal sessions transition to `TIMEOUT`.
- `generation_segment_cap` — a session that hits the cap transitions to
`COMPLETE` after the last segment ships.
## Message catalogue
Every JSON frame carries `{"type": <str>, ...}`. Pydantic models in
`protocol.py` are the source of truth; this table is the human-readable
view.
### Client → server
| `type` | Required fields | Purpose |
|---|---|---|
| `session_init_v2` | — | Opening frame. Carries preset, curated prompts, optional initial image, feature toggles, optional `continuation_state` to resume from a snapshot. |
| `segment_prompt_source` | `prompt` | Request the next segment using the supplied prompt; optional sampling overrides (`seed`, `num_inference_steps`, `guidance_scale`, `negative_prompt`). |
| `seed_prompts_updated` | `seed_prompts` | Replace the session's seed-prompt list; takes effect on the next segment. |
| `enhancement_updated` | `enabled` | Toggle prompt enhancement for subsequent segments. |
| `auto_extension_updated` | `enabled` | Toggle automatic per-segment prompt extension. |
| `loop_generation_updated` | `enabled` | Toggle loop-generation mode. |
| `generation_paused_updated` | `paused` | Pause/resume segment generation; queued requests defer. |
| `snapshot_state` | — | Request the current `ContinuationState` for export; server replies with `continuation_state_snapshot`. |
The opening frame must be `session_init_v2`. Any other first frame is
rejected with an `error` (code `invalid_message`) and the WebSocket is
closed.
### Server → client
| `type` | Carries | When emitted |
|---|---|---|
| `queue_status` | `position`, `queue_depth` | After `session_init_v2` accepted, before GPU binding. |
| `gpu_assigned` | GPU id, model id | Once a generator slot is bound. |
| `ltx2_stream_start` | session-level metadata | Once the session enters `ACTIVE`. |
| `ltx2_segment_start` | `segment_idx`, `prompt`, prompt source | When a `segment_prompt_source` request begins generation. |
| `step_complete` | `segment_idx`, denoise timings | After the segment's denoising loop finishes (before media emission). |
| `media_init` | `segment_idx`, mime, stream id | First frame of fMP4 output for the segment. |
| binary frame | fMP4 fragment bytes | Subsequent media chunks; the protocol enforces that `media_init` precedes any binary frames. |
| `media_segment_complete` | `segment_idx`, chunk count, byte count | Last media chunk for the segment. |
| `ltx2_segment_complete` | `segment_idx`, segment summary | Segment fully shipped; ready for the next `segment_prompt_source`. |
| `ltx2_stream_complete` | session summary | Session reached `generation_segment_cap` or client requested clean shutdown. |
| `session_timeout` | reason | Session hit `session_timeout_seconds`; immediately followed by close. |
| `continuation_state_snapshot` | `kind`, `payload` | Reply to `snapshot_state`. The payload is the same shape produced by `LTX2ContinuationState.to_continuation_state(...)`. |
| `error` | `code`, `message` | Any validation/runtime error. Non-fatal errors keep the connection open; fatal errors precede a `close`. |
## Continuation state
The session optionally accepts a `continuation_state` dict inside the
opening `session_init_v2` frame. When present, the server hydrates it
into a `ContinuationState(kind, payload)` envelope and feeds it as the
`request.state` on the first segment's `GenerationRequest` — letting a
client resume after a disconnect, migrate sessions across processes,
or replay a prior session.
After every segment, if the runtime returns a fresh state, the server
persists it to the `SessionStore` so a `snapshot_state` request can
export it. The store and serialization contracts live with the model
family (e.g. `fastvideo/pipelines/basic/ltx2/continuation.py` for LTX-2).
## Example flow
```
client server
────── ──────
WS /v1/stream ─────── connect ─────────────────────────►
◄────── (accept)
{"type": "session_init_v2",
"preset": "ltx2_two_stage",
"curated_prompts": ["a fox in snow", "the fox jumps"],
"initial_image": {...},
"stream_mode": "av_fmp4"} ─────────────────────────────►
(validate, queue, bind)
◄──── {"type": "queue_status",
"position": 0, "queue_depth": 0}
◄──── {"type": "gpu_assigned",
"gpu_id": 0, "model_id": "..."}
◄──── {"type": "ltx2_stream_start", ...}
{"type": "segment_prompt_source",
"prompt": "a fox in snow",
"source": "curated"} ───────────────────────────────────►
(run pipeline)
◄──── {"type": "ltx2_segment_start",
"segment_idx": 1, ...}
◄──── {"type": "step_complete",
"segment_idx": 1, "timings": {...}}
◄──── {"type": "media_init",
"segment_idx": 1,
"mime": "video/mp4", ...}
◄──── <binary fMP4 init segment>
◄──── <binary fMP4 fragment>
◄──── <binary fMP4 fragment>
◄──── {"type": "media_segment_complete",
"segment_idx": 1, "chunks": 12}
◄──── {"type": "ltx2_segment_complete",
"segment_idx": 1, ...}
{"type": "segment_prompt_source",
"prompt": "the fox jumps"} ─────────────────────────────►
(segment 2 …)
{"type": "snapshot_state"} ──────────────────────────────►
◄──── {"type": "continuation_state_snapshot",
"kind": "ltx2.v1",
"payload": {"schema_version": 1, ...}}
(close) ──────────────────────────────────────────────────►
(session → COMPLETE)
```
## Backward / forward compatibility
- Adding a new client message: append a Pydantic model to `protocol.py`
with a unique `type`; add the discriminator entry to `ClientMessage`;
add a row to the table above. Old clients that don't send the new
message remain compatible.
- Adding a new server message: emit only when a new feature flag is
enabled (or always emit, since clients ignore unknown types).
- Changing an existing message: bump the `type` (e.g. `session_init_v2`
→ `session_init_v3`) and accept both for one release cycle. Never
silently change field semantics under the same `type`.
+1 -1
View File
@@ -243,7 +243,7 @@ for step in range(start_step, max_steps):
```bash
# Install
uv pip install -e .[dev]
uv pip install -e ".[dev]"
# Run DMD2 distillation on Wan 2.1
torchrun --nproc_per_node=8 -m fastvideo.train.entrypoint.train \
+22
View File
@@ -86,3 +86,25 @@ sbatch examples/distill/Wan2.2-TI2V-5B-Diffusers/Data-free/distill_dmd_t2v_5B.sh
- Learning rate: 2e-5
- Training steps: 3000 (~12 hours)
- HSDP shard dim: 1
## 🧭 Note on `real_score_guidance_scale`
The teacher CFG used inside the DMD loss follows the DMD2 reference
implementation and uses the parameterization
```
x = x_cond + w * (x_cond - x_uncond)
```
rather than the Ho & Salimans form `x_uncond + w * (x_cond - x_uncond)`. The
two are mathematically equivalent up to a constant offset:
| `real_score_guidance_scale` (`w`) | Equivalent standard CFG (`w + 1`) | Output |
|-----------------------------------|-----------------------------------|-----------------------|
| `-1` | `0` | unconditional |
| `0` | `1` | conditional |
| `3.5` (default) | `4.5` | strong guidance |
So `real_score_guidance_scale` should be read as the **extra** guidance
strength added on top of the conditional prediction. When porting values
from a paper that uses the Ho & Salimans form, subtract 1.
+4 -4
View File
@@ -27,7 +27,7 @@ uv pip install fastvideo
conda create -n fastvideo python=3.12 -y
conda activate fastvideo
pip install fastvideo
uv pip install fastvideo
```
### From source
@@ -41,11 +41,11 @@ uv pip install -e .
uv pip install flash-attn --no-build-isolation -v
```
Alternative with Conda environment:
Alternative with Conda environment (still drives installs through `uv`):
```bash
pip install -e .
pip install flash-attn --no-build-isolation -v
uv pip install -e .
uv pip install flash-attn --no-build-isolation -v
```
## Hardware Requirements
+6 -4
View File
@@ -58,14 +58,16 @@ uv pip install flash-attn --no-build-isolation -v
#### With Conda environment (alternative)
`uv` works inside an active conda env too, so prefer `uv pip` for the actual install:
```bash
pip install fastvideo
uv pip install fastvideo
```
Also optionally install FlashAttention:
```bash
pip install flash-attn --no-build-isolation -v
uv pip install flash-attn --no-build-isolation -v
```
### Installation from Source
@@ -87,7 +89,7 @@ uv pip install -e .
Alternative with Conda environment:
```bash
pip install -e .
uv pip install -e .
```
### Optional Dependencies
@@ -101,7 +103,7 @@ uv pip install flash-attn --no-build-isolation -v
Alternative with Conda environment:
```bash
pip install flash-attn --no-build-isolation -v
uv pip install flash-attn --no-build-isolation -v
```
## Set up using Docker
+4 -2
View File
@@ -57,8 +57,10 @@ uv pip install fastvideo
#### With Conda environment (alternative)
`uv` works inside an active conda env too, so prefer `uv pip` for the actual install:
```bash
pip install fastvideo
uv pip install fastvideo
```
### Installation from Source
@@ -80,7 +82,7 @@ uv pip install -e .
Alternative with Conda environment:
```bash
pip install -e .
uv pip install -e .
```
## Development Environment Setup
+1 -1
View File
@@ -19,7 +19,7 @@
- Install MoGe:
```bash
pip install git+https://github.com/microsoft/MoGe.git
uv pip install git+https://github.com/microsoft/MoGe.git
```
- If you hit `ImportError: libGL.so.1` (common on Ubuntu/headless nodes), you can try installing OpenCV runtime libs:
+3 -3
View File
@@ -54,7 +54,7 @@ FASTVIDEO_ATTENTION_BACKEND=SAGE_ATTN python example.py
We recommend always installing [Flash Attention 2](https://github.com/Dao-AILab/flash-attention):
```bash
pip install flash-attn==2.7.4.post1 --no-build-isolation
uv pip install flash-attn==2.7.4.post1 --no-build-isolation
```
And if using a Hopper+ GPU (ie H100), installing [Flash Attention 3](https://github.com/Dao-AILab/flash-attention?tab=readme-ov-file#flashattention-3-beta-release) by compiling it from source (takes about 10 minutes for me):
@@ -63,7 +63,7 @@ And if using a Hopper+ GPU (ie H100), installing [Flash Attention 3](https://git
git clone https://github.com/Dao-AILab/flash-attention.git && cd flash-attention
cd hopper
pip install ninja
uv pip install ninja
python setup.py install
```
@@ -98,7 +98,7 @@ To use [SageAttention](https://github.com/thu-ml/SageAttention) 2.1.1, please co
```bash
git clone https://github.com/thu-ml/SageAttention.git
cd sageattention
python setup.py install # or pip install -e .
python setup.py install # or uv pip install -e .
```
### Sage Attention 3
@@ -4,7 +4,7 @@ These are end-to-end example scripts for distilling Wan2.1 T2V 1.3B model using
### 0. Make sure you have installed VSA
```bash
pip install vsa
uv pip install vsa
```
### 1. Download dataset:
@@ -4,7 +4,7 @@ These are end-to-end example scripts for distilling Wan2.2 TI2V 5B model DMD+VSA
### 0. Make sure you have installed VSA
```bash
pip install vsa
uv pip install vsa
```
### Data-free Distillation
@@ -4,7 +4,7 @@ These are end-to-end example scripts for distilling Wan2.2 TI2V 5B model DMD+VSA
### 0. Make sure you have installed VSA
```bash
pip install vsa
uv pip install vsa
```
### 1. Download dataset:
+1 -1
View File
@@ -7,7 +7,7 @@ and the GEN3C diffusion model.
Requirements:
1. Install MoGe:
pip install git+https://github.com/microsoft/MoGe.git
uv pip install git+https://github.com/microsoft/MoGe.git
If you hit `ImportError: libGL.so.1`, install:
sudo apt-get update && sudo apt-get install -y libgl1 libglib2.0-0 libsm6 libxext6 libxrender1
2. Download and convert weights:
@@ -0,0 +1,51 @@
# SPDX-License-Identifier: Apache-2.0
"""Minimal user-runnable example for the daVinci-MagiHuman base AV pipeline.
Produces an mp4 with both video (Wan 2.2 TI2V-5B VAE) and audio (Stable
Audio Open 1.0 VAE, first-class FastVideo port in
`fastvideo/models/vaes/oobleck.py`) muxed together via PyAV.
Prerequisites (one-off):
# Accept terms of use on the gated HF repos with your HF_TOKEN:
# - https://huggingface.co/google/t5gemma-9b-9b-ul2
# - https://huggingface.co/stabilityai/stable-audio-open-1.0
# All four cross-variant shared components (Wan 2.2 VAE, T5-Gemma
# encoder + tokenizer, Stable Audio VAE) are lazy-loaded from their
# canonical upstream HF repos on first build, so a single ~25 GB
# cache is shared across every MagiHuman variant.
The umbrella HF repo `FastVideo/MagiHuman-Diffusers` holds all four
variants (base / distill / sr_540p / sr_1080p) under sibling subfolders
and FastVideo will download just the requested subfolder. Local
conversion via `scripts/checkpoint_conversion/convert_magi_human_to_diffusers.py`
is also supported.
"""
from fastvideo import VideoGenerator
PROMPT = (
"A warm afternoon scene: a person sits on a park bench reading a book, "
"surrounded by softly swaying trees."
)
def main() -> None:
generator = VideoGenerator.from_pretrained(
"FastVideo/MagiHuman-Diffusers/base",
num_gpus=1,
)
output_path = "outputs_video/magi_human_basic/output_magi_human.mp4"
generator.generate_video(
prompt=PROMPT,
output_path=output_path,
save_video=True,
# Defaults pulled from the registered preset (magi_human_base):
# height=256, width=448, fps=25, num_inference_steps=32, seed=42.
# Override here only if you have a specific QA scenario.
)
generator.shutdown()
if __name__ == "__main__":
main()
@@ -0,0 +1,53 @@
# SPDX-License-Identifier: Apache-2.0
"""Minimal user-runnable example for the daVinci-MagiHuman DMD-2 distilled
text-to-AV pipeline.
Same arch as the base model (`basic_magi_human.py`) but with DMD-2 distilled
weights: 8 denoising steps, no classifier-free guidance. ~4x faster than
base at the same 256x480 resolution. Mirrors upstream
`daVinci-MagiHuman/example/distill/run_T2V.sh`.
Prerequisites (one-off):
# 1) Accept terms on the gated HF repos with your HF_TOKEN:
# - https://huggingface.co/google/t5gemma-9b-9b-ul2
# - https://huggingface.co/stabilityai/stable-audio-open-1.0
# Cross-variant shared components (Wan 2.2 VAE + T5-Gemma + Stable
# Audio VAE) are lazy-loaded from their canonical upstream HF repos
# and shared with the base variant cache.
# 2) Convert the distill subfolder of GAIR/daVinci-MagiHuman:
python scripts/checkpoint_conversion/convert_magi_human_to_diffusers.py \\
--source GAIR/daVinci-MagiHuman \\
--subfolder distill \\
--output converted_weights/magi_human_distill \\
--cast-bf16
# `--cast-bf16` is recommended (61 GB fp32 -> 30 GB bf16); the FV pipeline
# loads bf16 anyway, and the conversion keeps norms / RoPE bands fp32.
"""
from fastvideo import VideoGenerator
PROMPT = (
"A warm afternoon scene: a person sits on a park bench reading a book, "
"surrounded by softly swaying trees."
)
def main() -> None:
generator = VideoGenerator.from_pretrained(
"FastVideo/MagiHuman-Diffusers/distill",
num_gpus=1,
)
output_path = "outputs_video/magi_human_basic/output_magi_human_distill.mp4"
generator.generate_video(
prompt=PROMPT,
output_path=output_path,
save_video=True,
# Defaults pulled from the registered preset (magi_human_distill):
# height=256, width=480, fps=25, num_inference_steps=8, cfg=1, seed=42.
)
generator.shutdown()
if __name__ == "__main__":
main()
@@ -0,0 +1,34 @@
# SPDX-License-Identifier: Apache-2.0
"""Minimal daVinci-MagiHuman DMD-2 distilled text+image-to-AV example."""
from fastvideo import VideoGenerator
from fastvideo.pipelines.basic.magi_human.pipeline_configs import (
MagiHumanDistillI2VConfig,
)
PROMPT = (
"A cheerful saxophonist performs a short line with expressive facial "
"motion, natural head movement, and synchronized audio in a small jazz club."
)
IMAGE_PATH = "assets/images/saxophonist.jpg"
def main() -> None:
generator = VideoGenerator.from_pretrained(
"FastVideo/MagiHuman-Diffusers/distill",
num_gpus=1,
workload_type="i2v",
override_pipeline_cls_name="MagiHumanI2VPipeline",
pipeline_config=MagiHumanDistillI2VConfig(),
)
generator.generate_video(
prompt=PROMPT,
image_path=IMAGE_PATH,
output_path="outputs_video/magi_human_distill_ti2v/output_magi_human_distill_ti2v.mp4",
save_video=True,
)
generator.shutdown()
if __name__ == "__main__":
main()
@@ -0,0 +1,45 @@
# SPDX-License-Identifier: Apache-2.0
"""Run daVinci-MagiHuman SR-1080p text-to-AV in FastVideo.
Build the converted repo on large local storage, then symlink it into the
workspace:
python scripts/checkpoint_conversion/convert_magi_human_to_diffusers.py \
--source GAIR/daVinci-MagiHuman \
--subfolder base \
--sr-source GAIR/daVinci-MagiHuman \
--sr-subfolder 1080p_sr \
--output /raid/william5lin_converted_weights/magi_human_sr_1080p \
--cast-bf16
ln -s /raid/william5lin_converted_weights/magi_human_sr_1080p \
converted_weights/magi_human_sr_1080p
"""
from fastvideo import VideoGenerator
from fastvideo.pipelines.basic.magi_human.pipeline_configs import (
MagiHumanSR1080pConfig,
)
PROMPT = (
"A warm afternoon scene: a person sits on a park bench reading a book, "
"surrounded by softly swaying trees."
)
def main() -> None:
generator = VideoGenerator.from_pretrained(
"FastVideo/MagiHuman-Diffusers/sr_1080p",
num_gpus=1,
override_pipeline_cls_name="MagiHumanSR1080pPipeline",
pipeline_config=MagiHumanSR1080pConfig(),
)
generator.generate_video(
prompt=PROMPT,
output_path="outputs_video/magi_human_sr1080p/output_magi_human_sr1080p.mp4",
save_video=True,
)
generator.shutdown()
if __name__ == "__main__":
main()
@@ -0,0 +1,34 @@
# SPDX-License-Identifier: Apache-2.0
"""Run daVinci-MagiHuman SR-1080p text+image-to-AV in FastVideo."""
from fastvideo import VideoGenerator
from fastvideo.pipelines.basic.magi_human.pipeline_configs import (
MagiHumanSR1080pI2VConfig,
)
PROMPT = (
"A cheerful saxophonist performs a short line with expressive facial "
"motion, natural head movement, and synchronized audio in a small jazz club."
)
IMAGE_PATH = "assets/images/saxophonist.jpg"
def main() -> None:
generator = VideoGenerator.from_pretrained(
"FastVideo/MagiHuman-Diffusers/sr_1080p",
num_gpus=1,
workload_type="i2v",
override_pipeline_cls_name="MagiHumanSR1080pI2VPipeline",
pipeline_config=MagiHumanSR1080pI2VConfig(),
)
generator.generate_video(
prompt=PROMPT,
image_path=IMAGE_PATH,
output_path="outputs_video/magi_human_sr1080p_ti2v/output_magi_human_sr1080p_ti2v.mp4",
save_video=True,
)
generator.shutdown()
if __name__ == "__main__":
main()
@@ -0,0 +1,37 @@
# SPDX-License-Identifier: Apache-2.0
"""Run daVinci-MagiHuman SR-540p text-to-AV in FastVideo.
The converted repo must contain both ``transformer/`` (base DiT) and
``sr_transformer/`` (540p SR DiT). Build it with:
python scripts/checkpoint_conversion/convert_magi_human_to_diffusers.py \
--source GAIR/daVinci-MagiHuman \
--subfolder base \
--sr-source GAIR/daVinci-MagiHuman \
--sr-subfolder 540p_sr \
--output converted_weights/magi_human_sr_540p
"""
from fastvideo import VideoGenerator
PROMPT = (
"A warm afternoon scene: a person sits on a park bench reading a book, "
"surrounded by softly swaying trees."
)
def main() -> None:
generator = VideoGenerator.from_pretrained(
"FastVideo/MagiHuman-Diffusers/sr_540p",
num_gpus=1,
)
generator.generate_video(
prompt=PROMPT,
output_path="outputs_video/magi_human_sr540p/output_magi_human_sr540p.mp4",
save_video=True,
)
generator.shutdown()
if __name__ == "__main__":
main()
@@ -0,0 +1,34 @@
# SPDX-License-Identifier: Apache-2.0
"""Run daVinci-MagiHuman SR-540p text+image-to-AV in FastVideo."""
from fastvideo import VideoGenerator
from fastvideo.pipelines.basic.magi_human.pipeline_configs import (
MagiHumanSR540pI2VConfig,
)
PROMPT = (
"A cheerful saxophonist performs a short line with expressive facial "
"motion, natural head movement, and synchronized audio in a small jazz club."
)
IMAGE_PATH = "assets/images/saxophonist.jpg"
def main() -> None:
generator = VideoGenerator.from_pretrained(
"FastVideo/MagiHuman-Diffusers/sr_540p",
num_gpus=1,
workload_type="i2v",
override_pipeline_cls_name="MagiHumanSRI2VPipeline",
pipeline_config=MagiHumanSR540pI2VConfig(),
)
generator.generate_video(
prompt=PROMPT,
image_path=IMAGE_PATH,
output_path="outputs_video/magi_human_sr540p_ti2v/output_magi_human_sr540p_ti2v.mp4",
save_video=True,
)
generator.shutdown()
if __name__ == "__main__":
main()
@@ -0,0 +1,34 @@
# SPDX-License-Identifier: Apache-2.0
"""Minimal daVinci-MagiHuman base text+image-to-AV example."""
from fastvideo import VideoGenerator
from fastvideo.pipelines.basic.magi_human.pipeline_configs import (
MagiHumanBaseI2VConfig,
)
PROMPT = (
"A cheerful saxophonist performs a short line with expressive facial "
"motion, natural head movement, and synchronized audio in a small jazz club."
)
IMAGE_PATH = "assets/images/saxophonist.jpg"
def main() -> None:
generator = VideoGenerator.from_pretrained(
"FastVideo/MagiHuman-Diffusers/base",
num_gpus=1,
workload_type="i2v",
override_pipeline_cls_name="MagiHumanI2VPipeline",
pipeline_config=MagiHumanBaseI2VConfig(),
)
generator.generate_video(
prompt=PROMPT,
image_path=IMAGE_PATH,
output_path="outputs_video/magi_human_ti2v/output_magi_human_ti2v.mp4",
save_video=True,
)
generator.shutdown()
if __name__ == "__main__":
main()
@@ -0,0 +1,77 @@
# SPDX-License-Identifier: Apache-2.0
"""Stable Audio Open 1.0 — text-to-audio (baseline) example.
User story (game-audio designer, prototyping):
"I'm prototyping a level and I need 6 seconds of background
ambience — gentle wind, distant thunder, a hint of birdsong. I
don't want to dig through a sound library; I want to type what I
hear in my head and get a wav back. If it's wrong I'll iterate
on the prompt. This is the first stop."
User story (musician sketching ideas):
"I want to bounce a 30s lo-fi drum loop to use as a placeholder
bed while I build the rest of the track. Type prompt, get audio,
drop into the DAW. The actual production beat I'll record
myself, but I need *something* to write the chords against."
User story (researcher exploring the model):
"First time touching Stable Audio Open — what does it sound
like at default settings? This is the smallest amount of code
that goes from prompt to mp4."
How it works:
Pure text-to-audio (T2A). The pipeline runs:
T5 + NumberConditioner -> StableAudioDiT -> Oobleck VAE
via the `dpmpp-3m-sde` k-diffusion sampler. All components are
FastVideo-native — no diffusers / transformers model imports at
runtime (see REVIEW item 30). Mirrors upstream
`stable_audio_tools.inference.generation.generate_diffusion_cond`
bit-for-bit (~0.2% abs_mean drift on 25 steps).
Tunable knobs (the "creative dials"):
audio_end_in_s
1–6 — quick ideation (sub-10s wall clock at 100 steps)
10–30 — full musical phrase / loop length (the README example
uses 30s)
47.5 — model maximum (full sample_size = 2097152 / 44100 Hz)
num_inference_steps
25 — fast preview, occasional artifacts
100 — preset default (matches the HF model card)
250 — diminishing returns past here
guidance_scale
3 — looser, more variation per seed
7 — preset default; matches README
12+ — sharper but can sound "fried"
Prerequisites:
1. Accept the terms on https://huggingface.co/stabilityai/stable-audio-open-1.0
and export your HF token in the shell:
export HF_TOKEN=hf_...
2. Install optional inference deps (one-time):
uv pip install k_diffusion einops_exts alias_free_torch torchsde
"""
from fastvideo import VideoGenerator
PROMPT = "Lo-fi hip hop instrumental with vinyl crackle and gentle piano."
def main() -> None:
generator = VideoGenerator.from_pretrained(
"FastVideo/stable-audio-open-1.0-Diffusers",
num_gpus=1,
)
output_path = "outputs_audio/stable_audio_basic/output_stable_audio.wav"
generator.generate_video(
prompt=PROMPT,
output_path=output_path,
save_video=True,
# 6-second clip; the model max is ~47.5s.
audio_end_in_s=6.0,
# The registered preset gives 100 steps + CFG=7.0 by default;
# override num_inference_steps / guidance_scale here for QA.
)
generator.shutdown()
if __name__ == "__main__":
main()
@@ -0,0 +1,77 @@
# SPDX-License-Identifier: Apache-2.0
"""Stable Audio Open 1.0 — audio-to-audio variation example.
User story (musician, late at night):
"I generated this 12-second lo-fi loop earlier and I love the chord
progression and overall vibe, but the snare hit at 0:08 sounds wrong
and the rhythm feels stiff. I don't want to start over from scratch
and lose what's working — I want the model to keep the harmony and
mood but reroll the percussion + groove."
User story (sound designer, on a deadline):
"I have one good 'sword clang' SFX. The art director wants 8 sibling
variations that all feel like the same sword from different angles —
same metal, same weight, slightly different impact. I'd rather
refine my one good take than text-prompt my way through 50 misses."
Pass `init_audio=path/to/clip` (any wav/mp3/mp4/m4a/flac the standard
deps decode) and the model will use it as a starting point for the
text prompt instead of pure noise.
Picking `init_audio_strength` (0.0 to 1.0):
Higher = closer to the source clip. Lower = more transformation.
(Same convention as the "Input Audio Strength" slider in
Stability's commercial Stable Audio web UI, so values transfer
directly.)
| strength | what you get |
|----------|----------------------------------------------------|
| 1.00 | Output ≈ reference. No transformation. |
| 0.85 | Texture micro-variation only. |
| 0.70 | Light reroll, same instruments. |
| 0.60 | Default. Instrument identity is replaceable |
| | (cello can take over from piano on the same notes).|
| 0.50 | Heavy — only melody / chord progression survives. |
| 0.30 | Reference acts as a loose mood prompt. |
| 0.00 | Plain T2A — reference ignored. |
Rule of thumb by intent:
* "Fix one part of this clip" -> 0.75 .. 0.85
* "Same notes, different instrument" -> 0.55 .. 0.65
* "Same chord progression, new content" -> 0.40 .. 0.55
* "Use this as a loose mood prompt" -> 0.20 .. 0.35
If the reference timbre is bleeding through more than you want,
lower it; if the structure is gone, raise it.
Prerequisites: same as `basic_stable_audio.py`.
"""
from fastvideo import VideoGenerator
PROMPT = "Change the piano to a cello playing the same notes"
# Path to any audio-bearing file (wav, mp3, mp4, m4a, flac, ...).
# Set to `None` to skip A2A and run plain T2A.
INIT_AUDIO_PATH: str | None = None
# Reference fidelity in [0, 1] -- higher = closer to source.
INIT_AUDIO_STRENGTH = 0.6
def main() -> None:
generator = VideoGenerator.from_pretrained(
"FastVideo/stable-audio-open-1.0-Diffusers",
num_gpus=1,
)
generator.generate_video(
prompt=PROMPT,
output_path="outputs_audio/stable_audio_a2a/output_a2a.wav",
save_video=True,
audio_end_in_s=6.0,
init_audio=INIT_AUDIO_PATH,
init_audio_strength=INIT_AUDIO_STRENGTH,
)
generator.shutdown()
if __name__ == "__main__":
main()
@@ -0,0 +1,84 @@
# SPDX-License-Identifier: Apache-2.0
"""Stable Audio Open 1.0 — inpainting / outpainting (loop extension) example.
User story (loop extension — the killer app):
"I have a 6-second drum loop my client likes. They want it as
background bed for a 30-second ad. I need it to loop seamlessly,
but a hard cut every 6s sounds bad. Let me extend it to 30s,
keeping the first 6s exactly as-is and letting the model continue
the groove for the remaining 24s."
User story (audio repair):
"There's a microphone bump at 0:14 in this 30-second field
recording — really obvious in headphones. Mask out 0:13 to 0:15
and let the model regenerate plausible ambience that blends in.
Everything else stays exactly as I recorded it."
User story (transition smoothing):
"I have two 10-second clips I want to crossfade. Mask out a 1s
overlap region in the middle and let the model invent a coherent
transition between the two."
How it works (RePaint-style blending):
Stable Audio Open 1.0 wasn't trained as an inpainting model
(`model_type=diffusion_cond`, not `diffusion_cond_inpaint`), so we
can't use the upstream's mask-conditioned approach directly. We
use the RePaint trick instead, which works on any v-prediction
diffusion model:
1. Encode the reference clip into latent space.
2. At every denoising step `i`, replace the kept region of the
in-flight latent (where mask == 1) with the reference
re-noised to the next timestep's sigma. Only the unkept
region (mask == 0) is freely denoised.
3. After the loop, the kept region is exactly the reference;
the unkept region is freshly generated content.
This is approximate compared to a properly trained inpainting
checkpoint — the seam between kept/unkept can have slight EQ
discontinuity — but it works on the existing public model.
Tunable: the mask is a 1-D tensor in {0, 1} at the model's sample
rate. Conventions:
1.0 = keep this sample from the reference
0.0 = regenerate this sample
Prerequisites: same as `basic_stable_audio.py`.
"""
import os
from fastvideo import VideoGenerator
PROMPT = "Steady lo-fi hip hop drum loop with vinyl crackle."
# Required: path to the reference audio file (wav, mp3, mp4, m4a, flac,
# ...) you want to extend or repair. The pipeline raises if a mask is
# passed without a reference, so this must be a real path.
REFERENCE_AUDIO_PATH = "path/to/your/loop.wav"
KEEP_SECONDS = 6.0 # first KEEP_SECONDS preserved exactly
TOTAL_SECONDS = 12.0 # extend the loop to this duration
def main() -> None:
if not os.path.isfile(REFERENCE_AUDIO_PATH):
raise FileNotFoundError(
f"REFERENCE_AUDIO_PATH={REFERENCE_AUDIO_PATH!r} does not exist. "
"Edit this script to point at a real audio file (wav/mp3/mp4/"
"m4a/flac) before running.")
generator = VideoGenerator.from_pretrained(
"FastVideo/stable-audio-open-1.0-Diffusers",
num_gpus=1,
)
generator.generate_video(
prompt=PROMPT,
output_path="outputs_audio/stable_audio_inpaint/output_inpaint.wav",
save_video=True,
audio_end_in_s=TOTAL_SECONDS,
inpaint_audio=REFERENCE_AUDIO_PATH,
# Tuple form: keep first KEEP_SECONDS, regenerate the rest.
inpaint_mask=(KEEP_SECONDS, TOTAL_SECONDS),
)
generator.shutdown()
if __name__ == "__main__":
main()
@@ -0,0 +1,53 @@
# SPDX-License-Identifier: Apache-2.0
"""Stable Audio Open Small — fast / lightweight T2A example.
User story (interactive UI builder):
"I'm building a sound-design UI where the user types a prompt and
we want sub-2-second feedback so the experience feels like
autocomplete, not a render queue. The full Stable Audio Open 1.0
takes ~8s on a single GPU; the small variant takes a fraction of
that — quality is lower but completely usable for real-time
iteration."
User story (overnight batch jobs):
"I'm generating 10,000 short SFX variants for a procedural game.
Wall-clock matters more than per-clip polish — give me the small
model so I can fit the run in one night instead of a week."
How it works:
The small variant is a separate Stability AI checkpoint
(`stabilityai/stable-audio-open-small`) that ships the same Oobleck
VAE as the 1.0 base model but a smaller / faster DiT (`embed_dim=1024`,
`depth=16`, `qk_norm="ln"`) and only one duration conditioner
(`seconds_total`, no `seconds_start`). FastVideo loads from the
converted Diffusers-format repo `FastVideo/stable-audio-open-small-Diffusers`
via the standard component loader; per-variant arch fields come
from `transformer/config.json` and `conditioner/config.json`.
Prerequisites: same as `basic_stable_audio.py`. The converted repo is
public so no gated-access flow is required.
"""
from fastvideo import VideoGenerator
PROMPT = "Lo-fi hip hop instrumental with vinyl crackle and gentle piano."
def main() -> None:
generator = VideoGenerator.from_pretrained(
"FastVideo/stable-audio-open-small-Diffusers",
num_gpus=1,
)
output_path = "outputs_audio/stable_audio_small/output_stable_audio_small.wav"
generator.generate_video(
prompt=PROMPT,
output_path=output_path,
save_video=True,
# Small variant trains on a ~11.9s window — keep `audio_end_in_s`
# at or below that.
audio_end_in_s=6.0,
)
generator.shutdown()
if __name__ == "__main__":
main()
@@ -0,0 +1,73 @@
# Cosmos Predict2 2B T2V finetune config.
#
# Data must be preprocessed with Cosmos VAE + T5 text encoder
# into parquet format before training.
models:
student:
_target_: fastvideo.train.models.cosmos.CosmosModel
init_from: nvidia/Cosmos-Predict2-2B-Video2World
trainable: true
method:
_target_: fastvideo.train.methods.fine_tuning.finetune.FineTuneMethod
training:
distributed:
num_gpus: 8
sp_size: 1
tp_size: 1
hsdp_replicate_dim: 8
hsdp_shard_dim: 1
data:
data_path: data/cosmos_preprocessed
dataloader_num_workers: 4
train_batch_size: 1
training_cfg_rate: 0.0
seed: 1000
# Cosmos VAE: 4x temporal, 8x spatial compression.
# 93 frames -> 24 latent frames, 480x832 -> 60x104
num_latent_t: 24
num_height: 480
num_width: 832
num_frames: 93
optimizer:
learning_rate: 1.0e-5
betas: [0.9, 0.999]
weight_decay: 0.01
lr_scheduler: constant
lr_warmup_steps: 0
loop:
max_train_steps: 5000
gradient_accumulation_steps: 1
checkpoint:
output_dir: outputs/cosmos_finetune
training_state_checkpointing_steps: 500
checkpoints_total_limit: 3
resume_from_checkpoint: latest
tracker:
project_name: fastvideo_cosmos
run_name: cosmos_finetune
model:
enable_gradient_checkpointing_type: full
callbacks:
grad_clip:
_target_: fastvideo.train.callbacks.grad_clip.GradNormClipCallback
max_grad_norm: 1.0
validation:
_target_: fastvideo.train.callbacks.validation.ValidationCallback
pipeline_target: fastvideo.pipelines.basic.cosmos.cosmos_pipeline.Cosmos2VideoToWorldPipeline
dataset_file: data/cosmos_preprocessed/validation_prompts.json
every_steps: 100
sampling_steps: [50]
guidance_scale: 6.0
pipeline:
flow_shift: 1.0
@@ -0,0 +1,79 @@
# Cosmos-Predict2.5-2B Text-to-World overfitting test config.
#
# Overfits on a few short videos (480x832, 93 frames) to verify the
# Cosmos 2.5 training plugin works end-to-end.
#
# Preprocess data first:
# CUDA_VISIBLE_DEVICES=0 python fastvideo/pipelines/preprocess/preprocess_cosmos25_overfit.py
#
# Run:
# bash examples/train/run.sh examples/train/configs/overfit_cosmos25_t2w.yaml
models:
student:
_target_: fastvideo.train.models.cosmos.CosmosModel
init_from: KyleShao/Cosmos-Predict2.5-2B-Diffusers
trainable: true
enable_gradient_checkpointing_type: full
flow_shift: 1.0
method:
_target_: fastvideo.train.methods.fine_tuning.finetune.FineTuneMethod
training:
distributed:
num_gpus: 1
sp_size: 1
tp_size: 1
hsdp_replicate_dim: 1
hsdp_shard_dim: 1
data:
data_path: data/cosmos25_overfit_preprocessed
dataloader_num_workers: 0
train_batch_size: 1
training_cfg_rate: 0.0
seed: 42
num_latent_t: 24
num_height: 480
num_width: 832
num_frames: 93
optimizer:
learning_rate: 5.0e-5
betas: [0.9, 0.999]
weight_decay: 0.0
lr_scheduler: constant
lr_warmup_steps: 0
loop:
max_train_steps: 300
gradient_accumulation_steps: 1
checkpoint:
output_dir: outputs/cosmos25_overfit
training_state_checkpointing_steps: 50
checkpoints_total_limit: 2
tracker:
project_name: fastvideo_cosmos25
run_name: cosmos25_overfit
model:
precondition_outputs: false
enable_gradient_checkpointing_type: full
callbacks:
grad_clip:
_target_: fastvideo.train.callbacks.grad_clip.GradNormClipCallback
max_grad_norm: 1.0
validation:
_target_: fastvideo.train.callbacks.validation.ValidationCallback
pipeline_target: fastvideo.pipelines.basic.cosmos.cosmos2_5_pipeline.Cosmos2_5Pipeline
dataset_file: data/cosmos25_overfit_preprocessed/validation_prompts.json
every_steps: 150
sampling_steps: [35]
guidance_scale: 7.0
pipeline:
flow_shift: 1.0
+1 -1
View File
@@ -1,7 +1,7 @@
cmake_minimum_required(VERSION 3.26 FATAL_ERROR)
project(fastvideo-kernel LANGUAGES CXX)
# Prefer environment variable (used by CI or pip install git+repo_addr) if CMake var is not explicitly set.
# Prefer environment variable (used by CI or uv pip install git+repo_addr) if CMake var is not explicitly set.
if(NOT DEFINED GPU_BACKEND AND DEFINED ENV{GPU_BACKEND})
set(GPU_BACKEND "$ENV{GPU_BACKEND}")
endif()
+10 -1
View File
@@ -61,7 +61,16 @@ has_cmake_arg() {
}
detect_with_torch() {
uv run --active --no-project python -c "import torch
# Prefer the active venv's python directly over `uv run --active --no-project`,
# which on some uv versions provisions its own interpreter and misses packages
# installed into VIRTUAL_ENV.
local py
if [[ -n "${VIRTUAL_ENV:-}" && -x "${VIRTUAL_ENV}/bin/python" ]]; then
py="${VIRTUAL_ENV}/bin/python"
else
py="$(command -v python3 || command -v python)"
fi
"${py}" -c "import torch
if not torch.cuda.is_available():
raise RuntimeError('torch.cuda.is_available() is false')
mj, mn = torch.cuda.get_device_capability(0)
@@ -5,6 +5,11 @@ from fastvideo_kernel.ops import (
video_sparse_attn,
)
from fastvideo_kernel.block_sparse_attn import (
block_sparse_attn,
block_sparse_attn_from_indices,
)
from fastvideo_kernel.vmoba import (
moba_attn_varlen,
process_moba_input,
@@ -22,6 +27,8 @@ from fastvideo_kernel.turbodiffusion_ops import (
__all__ = [
"sliding_tile_attention",
"video_sparse_attn",
"block_sparse_attn",
"block_sparse_attn_from_indices",
"moba_attn_varlen",
"process_moba_input",
"process_moba_output",
@@ -1,3 +1,5 @@
"""Autograd-enabled block-sparse attention. Index-native ops with a bool-mask compat shim."""
from __future__ import annotations
import os
@@ -6,6 +8,11 @@ from typing import Tuple
import torch
# ---------------------------------------------------------------------------
# Backend selection helpers
# ---------------------------------------------------------------------------
def _get_sm90_ops():
try:
from fastvideo_kernel._C import fastvideo_kernel_ops # type: ignore
@@ -25,38 +32,66 @@ def _is_sm90() -> bool:
def _force_triton() -> bool:
# Force Triton even on SM90 and even if the compiled extension is available.
# Useful for CI / debugging / parity testing.
return os.environ.get("FASTVIDEO_KERNEL_VSA_FORCE_TRITON", "0") == "1"
def _map_to_index(block_map: torch.Tensor) -> Tuple[torch.Tensor, torch.Tensor]:
"""
Preferred map->index conversion used by the wrapper.
# ---------------------------------------------------------------------------
# Index helpers
# ---------------------------------------------------------------------------
This wrapper **requires** the Triton implementation.
If Triton (or the Triton map_to_index module) is not available, it raises.
"""
def _map_to_index(block_map: torch.Tensor) -> Tuple[torch.Tensor, torch.Tensor]:
"""Compact a bool block_map to (q2k_idx, q2k_num). Legacy path only."""
if block_map.dim() == 3:
block_map = block_map.unsqueeze(0)
if block_map.dim() != 4:
raise ValueError(f"block_map must be [B,H,Q,KV] (or [H,Q,KV]), got shape={tuple(block_map.shape)}")
raise ValueError(
f"block_map must be [B,H,Q,KV] (or [H,Q,KV]), "
f"got shape={tuple(block_map.shape)}"
)
if block_map.dtype != torch.bool:
block_map = block_map.to(torch.bool)
if not block_map.is_cuda:
raise RuntimeError("block_map must be a CUDA tensor (Triton map_to_index required).")
raise RuntimeError(
"block_map must be a CUDA tensor (Triton map_to_index required)."
)
try:
from fastvideo_kernel.triton_kernels.index import map_to_index as triton_map_to_index # local import
except Exception as e:
from fastvideo_kernel.triton_kernels.index import map_to_index as triton_map_to_index
except Exception as e: # pragma: no cover - environment issue
raise ImportError(
"Triton map_to_index is required but not available. "
"Ensure Triton is installed and fastvideo_kernel.triton_kernels.index is importable."
"Ensure Triton is installed and "
"fastvideo_kernel.triton_kernels.index is importable."
) from e
return triton_map_to_index(block_map)
def _invert_indices_for_backward(
q2k_idx: torch.Tensor,
q2k_num: torch.Tensor,
num_kv_blocks: int,
) -> Tuple[torch.Tensor, torch.Tensor]:
from fastvideo_kernel.triton_kernels.index import invert_indices
return invert_indices(q2k_idx, q2k_num, num_kv_blocks=num_kv_blocks)
def _as_int32_contig(t: torch.Tensor, name: str) -> torch.Tensor:
"""Return `t` as a contiguous int32 tensor, raising a clear error on CPU input."""
if not t.is_cuda:
raise RuntimeError(f"{name} must be a CUDA tensor, got device={t.device}")
if t.dtype != torch.int32:
t = t.to(torch.int32)
if not t.is_contiguous():
t = t.contiguous()
return t
# ---------------------------------------------------------------------------
# Triton backend custom ops (index-native)
# ---------------------------------------------------------------------------
@torch.library.custom_op(
"fastvideo_kernel::block_sparse_attn_triton",
mutates_args=(),
@@ -66,34 +101,40 @@ def block_sparse_attn_triton(
q: torch.Tensor,
k: torch.Tensor,
v: torch.Tensor,
block_map: torch.Tensor,
q2k_idx: torch.Tensor,
q2k_num: torch.Tensor,
variable_block_sizes: torch.Tensor,
) -> Tuple[torch.Tensor, torch.Tensor]:
q = q.contiguous()
k = k.contiguous()
v = v.contiguous()
block_map = block_map.to(torch.bool)
q2k_idx, q2k_num = _map_to_index(block_map)
from fastvideo_kernel.triton_kernels.block_sparse_attn_triton import ( # local import
from fastvideo_kernel.triton_kernels.block_sparse_attn_triton import (
triton_block_sparse_attn_forward,
)
o, M = triton_block_sparse_attn_forward(q, k, v, q2k_idx, q2k_num, variable_block_sizes)
o, M = triton_block_sparse_attn_forward(
q.contiguous(),
k.contiguous(),
v.contiguous(),
q2k_idx,
q2k_num,
variable_block_sizes,
)
return o, M
@torch.library.register_fake("fastvideo_kernel::block_sparse_attn_triton")
def _block_sparse_attn_triton_fake(
q: torch.Tensor,
k: torch.Tensor,
v: torch.Tensor,
block_map: torch.Tensor,
q2k_idx: torch.Tensor,
q2k_num: torch.Tensor,
variable_block_sizes: torch.Tensor,
) -> Tuple[torch.Tensor, torch.Tensor]:
o = torch.empty_like(q)
M = torch.empty((q.shape[0], q.shape[1], q.shape[2]), device=q.device, dtype=torch.float32)
M = torch.empty(
(q.shape[0], q.shape[1], q.shape[2]),
device=q.device,
dtype=torch.float32,
)
return o, M
@@ -109,20 +150,32 @@ def block_sparse_attn_backward_triton(
v: torch.Tensor,
o: torch.Tensor,
M: torch.Tensor,
block_map: torch.Tensor,
q2k_idx: torch.Tensor,
q2k_num: torch.Tensor,
variable_block_sizes: torch.Tensor,
) -> Tuple[torch.Tensor, torch.Tensor, torch.Tensor]:
grad_output = grad_output.contiguous()
block_map = block_map.to(torch.bool)
q2k_idx, q2k_num = _map_to_index(block_map)
k2q_idx, k2q_num = _map_to_index(block_map.transpose(-1, -2).contiguous())
from fastvideo_kernel.triton_kernels.block_sparse_attn_triton import ( # local import
from fastvideo_kernel.triton_kernels.block_sparse_attn_triton import (
triton_block_sparse_attn_backward,
)
num_kv_blocks = int(variable_block_sizes.numel())
k2q_idx, k2q_num = _invert_indices_for_backward(
q2k_idx, q2k_num, num_kv_blocks
)
# q/k/v are saved from the user-facing inputs and may be non-contiguous;
# o/M are kernel outputs so are already contiguous.
dq, dk, dv = triton_block_sparse_attn_backward(
grad_output, q, k, v, o, M, q2k_idx, q2k_num, k2q_idx, k2q_num, variable_block_sizes
grad_output.contiguous(),
q.contiguous(),
k.contiguous(),
v.contiguous(),
o,
M,
q2k_idx,
q2k_num,
k2q_idx,
k2q_num,
variable_block_sizes,
)
return dq, dk, dv
@@ -135,7 +188,8 @@ def _block_sparse_attn_backward_triton_fake(
v: torch.Tensor,
o: torch.Tensor,
M: torch.Tensor,
block_map: torch.Tensor,
q2k_idx: torch.Tensor,
q2k_num: torch.Tensor,
variable_block_sizes: torch.Tensor,
) -> Tuple[torch.Tensor, torch.Tensor, torch.Tensor]:
dq = torch.empty_like(q)
@@ -144,19 +198,28 @@ def _block_sparse_attn_backward_triton_fake(
return dq, dk, dv
def _backward_triton(ctx, grad_o, grad_M):
q, k, v, o, M, block_map, variable_block_sizes = ctx.saved_tensors
dq, dk, dv = block_sparse_attn_backward_triton(grad_o, q, k, v, o, M, block_map, variable_block_sizes)
return dq, dk, dv, None, None
def _setup_context_triton(ctx, inputs, output):
q, k, v, block_map, variable_block_sizes = inputs
q, k, v, q2k_idx, q2k_num, variable_block_sizes = inputs
o, M = output
ctx.save_for_backward(q, k, v, o, M, block_map, variable_block_sizes)
ctx.save_for_backward(q, k, v, o, M, q2k_idx, q2k_num, variable_block_sizes)
block_sparse_attn_triton.register_autograd(_backward_triton, setup_context=_setup_context_triton)
def _backward_triton(ctx, grad_o, grad_M):
q, k, v, o, M, q2k_idx, q2k_num, variable_block_sizes = ctx.saved_tensors
dq, dk, dv = block_sparse_attn_backward_triton(
grad_o, q, k, v, o, M, q2k_idx, q2k_num, variable_block_sizes
)
return dq, dk, dv, None, None, None
block_sparse_attn_triton.register_autograd(
_backward_triton, setup_context=_setup_context_triton
)
# ---------------------------------------------------------------------------
# SM90 backend custom ops (index-native)
# ---------------------------------------------------------------------------
@torch.library.custom_op(
@@ -168,21 +231,21 @@ def block_sparse_attn_sm90(
q_padded: torch.Tensor,
k_padded: torch.Tensor,
v_padded: torch.Tensor,
block_map: torch.Tensor,
q2k_idx: torch.Tensor,
q2k_num: torch.Tensor,
variable_block_sizes: torch.Tensor,
) -> Tuple[torch.Tensor, torch.Tensor]:
block_sparse_fwd, _ = _get_sm90_ops()
if block_sparse_fwd is None:
raise ImportError("fastvideo_kernel_ops.block_sparse_fwd is not available")
q_padded = q_padded.contiguous()
k_padded = k_padded.contiguous()
v_padded = v_padded.contiguous()
block_map = block_map.to(torch.bool)
q2k_idx, q2k_num = _map_to_index(block_map)
o_padded, lse_padded = block_sparse_fwd(
q_padded, k_padded, v_padded, q2k_idx, q2k_num, variable_block_sizes.int()
q_padded.contiguous(),
k_padded.contiguous(),
v_padded.contiguous(),
q2k_idx,
q2k_num,
variable_block_sizes,
)
return o_padded, lse_padded
@@ -192,11 +255,16 @@ def _block_sparse_attn_sm90_fake(
q_padded: torch.Tensor,
k_padded: torch.Tensor,
v_padded: torch.Tensor,
block_map: torch.Tensor,
q2k_idx: torch.Tensor,
q2k_num: torch.Tensor,
variable_block_sizes: torch.Tensor,
) -> Tuple[torch.Tensor, torch.Tensor]:
o = torch.empty_like(q_padded)
lse = torch.empty((q_padded.shape[0], q_padded.shape[1], q_padded.shape[2], 1), device=q_padded.device, dtype=torch.float32)
lse = torch.empty(
(q_padded.shape[0], q_padded.shape[1], q_padded.shape[2], 1),
device=q_padded.device,
dtype=torch.float32,
)
return o, lse
@@ -212,30 +280,34 @@ def block_sparse_attn_backward_sm90(
v_padded: torch.Tensor,
o_padded: torch.Tensor,
lse_padded: torch.Tensor,
block_map: torch.Tensor,
q2k_idx: torch.Tensor,
q2k_num: torch.Tensor,
variable_block_sizes: torch.Tensor,
) -> Tuple[torch.Tensor, torch.Tensor, torch.Tensor]:
_, block_sparse_bwd = _get_sm90_ops()
if block_sparse_bwd is None:
raise ImportError("fastvideo_kernel_ops.block_sparse_bwd is not available")
grad_output_padded = grad_output_padded.contiguous()
block_map = block_map.to(torch.bool)
k2q_idx, k2q_num = _map_to_index(block_map.transpose(-1, -2).contiguous())
num_kv_blocks = int(variable_block_sizes.numel())
k2q_idx, k2q_num = _invert_indices_for_backward(
q2k_idx, q2k_num, num_kv_blocks
)
# q/k/v are saved from user-facing inputs; o/lse are kernel outputs.
dq, dk, dv = block_sparse_bwd(
q_padded,
k_padded,
v_padded,
q_padded.contiguous(),
k_padded.contiguous(),
v_padded.contiguous(),
o_padded,
lse_padded,
grad_output_padded,
grad_output_padded.contiguous(),
k2q_idx,
k2q_num,
variable_block_sizes.int(),
variable_block_sizes,
)
# C++ kernel returns fp32 grads; cast back to match PyTorch convention if needed
return dq.to(grad_output_padded.dtype), dk.to(grad_output_padded.dtype), dv.to(grad_output_padded.dtype)
# C++ kernel returns fp32 grads; cast back to the input dtype.
out_dtype = grad_output_padded.dtype
return dq.to(out_dtype), dk.to(out_dtype), dv.to(out_dtype)
@torch.library.register_fake("fastvideo_kernel::block_sparse_attn_backward_sm90")
@@ -246,7 +318,8 @@ def _block_sparse_attn_backward_sm90_fake(
v_padded: torch.Tensor,
o_padded: torch.Tensor,
lse_padded: torch.Tensor,
block_map: torch.Tensor,
q2k_idx: torch.Tensor,
q2k_num: torch.Tensor,
variable_block_sizes: torch.Tensor,
) -> Tuple[torch.Tensor, torch.Tensor, torch.Tensor]:
dq = torch.empty_like(q_padded)
@@ -255,21 +328,57 @@ def _block_sparse_attn_backward_sm90_fake(
return dq, dk, dv
def _backward_sm90(ctx, grad_o, grad_lse):
q, k, v, o, lse, block_map, variable_block_sizes = ctx.saved_tensors
dq, dk, dv = block_sparse_attn_backward_sm90(
grad_o, q, k, v, o, lse, block_map, variable_block_sizes
)
return dq, dk, dv, None, None
def _setup_context_sm90(ctx, inputs, output):
q, k, v, block_map, variable_block_sizes = inputs
q, k, v, q2k_idx, q2k_num, variable_block_sizes = inputs
o, lse = output
ctx.save_for_backward(q, k, v, o, lse, block_map, variable_block_sizes)
ctx.save_for_backward(q, k, v, o, lse, q2k_idx, q2k_num, variable_block_sizes)
block_sparse_attn_sm90.register_autograd(_backward_sm90, setup_context=_setup_context_sm90)
def _backward_sm90(ctx, grad_o, grad_lse):
q, k, v, o, lse, q2k_idx, q2k_num, variable_block_sizes = ctx.saved_tensors
dq, dk, dv = block_sparse_attn_backward_sm90(
grad_o, q, k, v, o, lse, q2k_idx, q2k_num, variable_block_sizes
)
return dq, dk, dv, None, None, None
block_sparse_attn_sm90.register_autograd(
_backward_sm90, setup_context=_setup_context_sm90
)
# ---------------------------------------------------------------------------
# Public API
# ---------------------------------------------------------------------------
def block_sparse_attn_from_indices(
q: torch.Tensor,
k: torch.Tensor,
v: torch.Tensor,
q2k_idx: torch.Tensor,
q2k_num: torch.Tensor,
variable_block_sizes: torch.Tensor,
) -> Tuple[torch.Tensor, torch.Tensor]:
"""Block-sparse attention with autograd, taking compact per-row KV indices."""
# Normalize index tensors once at the public boundary so the custom ops
# and their fakes can assume int32/contiguous. No-op on well-formed input.
q2k_idx = _as_int32_contig(q2k_idx, "q2k_idx")
q2k_num = _as_int32_contig(q2k_num, "q2k_num")
variable_block_sizes = _as_int32_contig(variable_block_sizes, "variable_block_sizes")
block_sparse_fwd, block_sparse_bwd = _get_sm90_ops()
use_sm90 = (
(not _force_triton())
and _is_sm90()
and block_sparse_fwd is not None
and block_sparse_bwd is not None
)
if use_sm90:
return block_sparse_attn_sm90(q, k, v, q2k_idx, q2k_num, variable_block_sizes)
# Triton path: supports q_seq_len != kv_seq_len as long as both are padded
# to a multiple of the block size (64 tokens).
return block_sparse_attn_triton(q, k, v, q2k_idx, q2k_num, variable_block_sizes)
def block_sparse_attn(
@@ -279,16 +388,8 @@ def block_sparse_attn(
block_map: torch.Tensor,
variable_block_sizes: torch.Tensor,
) -> Tuple[torch.Tensor, torch.Tensor]:
"""
Unified block-sparse attention op with autograd support.
- On SM90 with compiled extension present: uses fastvideo_kernel_ops.block_sparse_fwd/bwd.
- Otherwise: uses Triton implementation (requires q/k/v to have same padded length today).
"""
block_sparse_fwd, block_sparse_bwd = _get_sm90_ops()
if (not _force_triton()) and _is_sm90() and (block_sparse_fwd is not None) and (block_sparse_bwd is not None):
return block_sparse_attn_sm90(q, k, v, block_map, variable_block_sizes)
# Triton path: supports q_seq_len != kv_seq_len as long as both are padded
# to a multiple of the block size (64 tokens).
return block_sparse_attn_triton(q, k, v, block_map, variable_block_sizes)
"""Bool-mask compat wrapper; prefer block_sparse_attn_from_indices."""
q2k_idx, q2k_num = _map_to_index(block_map)
return block_sparse_attn_from_indices(
q, k, v, q2k_idx, q2k_num, variable_block_sizes
)
@@ -1,6 +1,6 @@
import math
import torch
from .block_sparse_attn import block_sparse_attn
from .block_sparse_attn import block_sparse_attn, block_sparse_attn_from_indices
from .triton_kernels.st_attn_triton import sliding_tile_attention_triton
# Try to load the C++ extension
@@ -125,13 +125,18 @@ def video_sparse_attn(
out_c = out_c.repeat(1, 1, 1, block_elements,
1).view(batch, heads, q_seq_len, dim)
# Sparse branch
# Sparse branch: feed top-k indices directly, skipping the bool-mask round-trip.
topk_idx = torch.topk(scores, topk, dim=-1).indices
mask = torch.zeros_like(scores,
dtype=torch.bool).scatter_(-1, topk_idx, True)
# out_s = block_sparse_attn(q, k, v, mask, variable_block_sizes)[0]
out_s = block_sparse_attn(q, k, v, mask, variable_block_sizes)[0]
q2k_idx = topk_idx.to(torch.int32).contiguous()
q2k_num = torch.full(
(batch, heads, q_num_blocks),
topk,
dtype=torch.int32,
device=q.device,
)
out_s = block_sparse_attn_from_indices(
q, k, v, q2k_idx, q2k_num, variable_block_sizes
)[0]
if compress_attn_weight is not None:
return out_c * compress_attn_weight + out_s
@@ -1,9 +1,10 @@
## pytorch sdpa version of block sparse ##
from typing import Tuple
import triton
import triton.language as tl
import torch
@triton.jit
def topk_index_to_map_kernel(
map_ptr,
@@ -153,3 +154,114 @@ def map_to_index(block_map: torch.Tensor):
)
return index, index_num
@triton.jit
def _invert_indices_kernel(
q2k_idx_ptr,
q2k_num_ptr,
k2q_idx_ptr,
k2q_num_ptr,
q2k_idx_b, q2k_idx_h, q2k_idx_q, q2k_idx_k,
q2k_num_b, q2k_num_h, q2k_num_q,
k2q_idx_b, k2q_idx_h, k2q_idx_k, k2q_idx_q,
k2q_num_b, k2q_num_h, k2q_num_k,
MAX_KV_PER_Q: tl.constexpr,
):
# One program per (b, h, q): reserve a slot in k2q via atomicAdd, write q.
pid_b = tl.program_id(0)
pid_h = tl.program_id(1)
pid_q = tl.program_id(2)
n = tl.load(
q2k_num_ptr
+ pid_b * q2k_num_b
+ pid_h * q2k_num_h
+ pid_q * q2k_num_q
)
q2k_row = (
q2k_idx_ptr
+ pid_b * q2k_idx_b
+ pid_h * q2k_idx_h
+ pid_q * q2k_idx_q
)
for i in tl.range(0, MAX_KV_PER_Q):
if i < n:
kv = tl.load(q2k_row + i * q2k_idx_k)
count_ptr = (
k2q_num_ptr
+ pid_b * k2q_num_b
+ pid_h * k2q_num_h
+ kv * k2q_num_k
)
pos = tl.atomic_add(count_ptr, 1)
tl.store(
k2q_idx_ptr
+ pid_b * k2q_idx_b
+ pid_h * k2q_idx_h
+ kv * k2q_idx_k
+ pos * k2q_idx_q,
pid_q,
)
def invert_indices(
q2k_idx: torch.Tensor,
q2k_num: torch.Tensor,
num_kv_blocks: int,
) -> Tuple[torch.Tensor, torch.Tensor]:
"""Transpose a Q->KV index list into a K->Q one via atomic compaction (GPU)."""
if q2k_idx.dim() != 4:
raise ValueError(
f"q2k_idx must be [B, H, Nq, Mk], got shape={tuple(q2k_idx.shape)}"
)
if q2k_num.dim() != 3:
raise ValueError(
f"q2k_num must be [B, H, Nq], got shape={tuple(q2k_num.shape)}"
)
if not q2k_idx.is_cuda or not q2k_num.is_cuda:
raise RuntimeError("invert_indices requires CUDA tensors.")
B, H, Nq, Mk = q2k_idx.shape
if q2k_num.shape != (B, H, Nq):
raise ValueError(
f"q2k_num shape {tuple(q2k_num.shape)} does not match q2k_idx "
f"[B, H, Nq] = {(B, H, Nq)}"
)
q2k_idx = q2k_idx.contiguous()
q2k_num = q2k_num.contiguous()
if q2k_idx.dtype != torch.int32:
q2k_idx = q2k_idx.to(torch.int32)
if q2k_num.dtype != torch.int32:
q2k_num = q2k_num.to(torch.int32)
# Any KV block is attended by at most Nq Q blocks (one per Q row), so
# `Nq` is a tight upper bound on the compacted K->Q slots.
k2q_idx = torch.empty(
(B, H, num_kv_blocks, Nq),
dtype=torch.int32,
device=q2k_idx.device,
)
k2q_num = torch.zeros(
(B, H, num_kv_blocks),
dtype=torch.int32,
device=q2k_idx.device,
)
grid = (B, H, Nq)
_invert_indices_kernel[grid](
q2k_idx,
q2k_num,
k2q_idx,
k2q_num,
q2k_idx.stride(0), q2k_idx.stride(1), q2k_idx.stride(2), q2k_idx.stride(3),
q2k_num.stride(0), q2k_num.stride(1), q2k_num.stride(2),
k2q_idx.stride(0), k2q_idx.stride(1), k2q_idx.stride(2), k2q_idx.stride(3),
k2q_num.stride(0), k2q_num.stride(1), k2q_num.stride(2),
MAX_KV_PER_Q=Mk,
)
return k2q_idx, k2q_num
@@ -11,7 +11,7 @@ except ImportError:
def _unsupported(*args, **kwargs):
raise ImportError(
"flash-attn is not installed. Please install it, e.g., `pip install flash-attn`."
"flash-attn is not installed. Please install it, e.g., `uv pip install flash-attn`."
)
_flash_attn_varlen_forward = _unsupported
+67
View File
@@ -0,0 +1,67 @@
# `fastvideo/` — Core Package
**Generated:** 2026-05-02
Inference + training framework for video DiTs. Public API entry: `from fastvideo import VideoGenerator, PipelineConfig, SamplingParam`.
## Public Surface (`__init__.py`)
```python
VideoGenerator # entrypoints/video_generator.py — high-level inference handle
PipelineConfig # configs/pipelines/base.py — pipeline wiring dataclass
SamplingParam # api/sampling_param.py — runtime sampling knobs
```
CLI entry: `fastvideo` script → `entrypoints/cli/main.py` (subcommands: `generate`, `serve`, `bench`).
## Layout
```
fastvideo/
├── api/ # Schema + presets for the OpenAI-compatible serving layer
├── attention/ # Backends + selector (FlashAttn / SageAttn / SDPA / VSA / VMoBA / SLA)
├── configs/ # Per-model arch configs + per-pipeline configs (registry-driven)
├── dataset/ # Dataloaders (pre-commit excluded — minimal lint surface)
├── distributed/ # SP/TP groups, device communicators, init helpers
├── entrypoints/ # cli/, openai/, streaming/, video_generator.py
├── hooks/ # Runtime hook system for pipelines
├── layers/ # Tensor-parallel linears + attention wrappers (port targets)
├── models/ # DiT / VAE / encoder / scheduler / loader (pre-commit excluded)
├── pipelines/ # basic/<model>/, preprocess/, stages/, training/
├── platforms/ # CUDA/ROCm capability + AttentionBackendEnum
├── third_party/ # Vendored externals (lint excluded; do not reformat)
├── train/ # NEW modular trainer — methods × models × callbacks
├── training/ # LEGACY monolithic *_training/distillation_pipeline.py
├── worker/ # Multi-process / Ray executors
├── workflow/ # Preprocessing workflow base class
├── registry.py # Pipeline-config + model-class lookup (canonical)
├── envs.py # Env-var declarations
├── fastvideo_args.py# Runtime arg dataclass passed through pipelines
└── utils.py # FlexibleArgumentParser, qualname resolver, etc.
```
## Where to Look
| Task | Location |
|------|----------|
| Add a new pipeline class | `pipelines/basic/<model>/` + `configs/pipelines/<model>.py` + register in `registry.py` |
| Add a new model component | `models/<role>/<model>.py` + `configs/models/<role>/<model>.py` |
| Wire an existing model into a new pipeline | `pipelines/basic/<model>/presets.py` + reuse stages from `pipelines/stages/` |
| Add a converter | `scripts/checkpoint_conversion/<model>_to_*.py` (separate dir, separate AGENTS.md) |
| Add an attention backend | `attention/backends/<name>.py` + register in selector |
| Add a runtime CLI flag | `fastvideo_args.py` (avoid `argparse` ad-hoc inside stages) |
## Conventions Specific Here
- `PipelineStage` subclasses (`pipelines/stages/`) own one verb each (encode, schedule, denoise, decode). Compose, don't fork.
- Every pipeline reads from a `PipelineConfig` subclass and a `SamplingParam`. Never read raw env vars inside a stage — go through `fastvideo.envs`.
- Logger setup: `from fastvideo.logger import init_logger; logger = init_logger(__name__)`. Do not call `logging.getLogger` directly.
- Imports between `train/` and `training/` are **forbidden** — they are independent stacks.
## Pre-Commit Exclusions (do not assume linted)
These dirs are listed in `.pre-commit-config.yaml` `exclude`:
- `fastvideo/third_party/`, `fastvideo/dataset/`, `fastvideo/models/`
Editing files there will NOT trigger yapf/ruff/mypy/codespell. Format manually if a sibling file shows clear style; do not introduce new violations.
+118 -9
View File
@@ -16,6 +16,8 @@ from fastvideo.api.request_metadata import (
reset_tracking_roots,
)
from fastvideo.api.schema import (
CompileConfig,
ContinuationState,
GenerationRequest,
GeneratorConfig,
InputConfig,
@@ -25,6 +27,10 @@ from fastvideo.api.schema import (
)
from fastvideo.api.sampling_param import SamplingParam
from fastvideo.fastvideo_args import FastVideoArgs
from fastvideo.pipelines.basic.ltx2.stage_overrides import (
refine_preset_override_fields,
refine_stage_override_fields,
)
from fastvideo.utils import shallow_asdict
_INPUT_FIELD_NAMES = {field.name for field in fields(InputConfig)}
@@ -38,6 +44,10 @@ _LEGACY_REQUEST_ALIASES = {
_REQUEST_PIPELINE_OVERRIDE_FIELDS = frozenset({
"embedded_cfg_scale",
})
# torch.compile kwargs that map to first-class CompileConfig fields.
_COMPILE_TYPED_KEYS = ("backend", "fullgraph", "mode", "dynamic")
# LTX-2 refine flat kwargs (init + per-request) known to FastVideoArgs.
_LTX2_REFINE_FLAT_KEYS = (refine_preset_override_fields() | refine_stage_override_fields())
def normalize_generator_config(config: GeneratorConfig | Mapping[str, Any], ) -> GeneratorConfig:
@@ -80,6 +90,8 @@ def legacy_from_pretrained_to_config(
components: dict[str, Any] = {}
quantization: dict[str, Any] = {}
experimental: dict[str, Any] = {}
preset_overrides: dict[str, Any] = {}
preset_refine: dict[str, Any] = {}
for key, value in kwargs.items():
if key == "revision":
@@ -106,8 +118,33 @@ def legacy_from_pretrained_to_config(
offload["pin_cpu_memory"] = value
elif key == "enable_torch_compile":
compile_config["enabled"] = value
elif key == "enable_torch_compile_text_encoder":
compile_config["text_encoder_enabled"] = value
elif key == "torch_compile_kwargs":
compile_config["kwargs"] = deepcopy(value)
remaining: dict[str, Any] = (dict(deepcopy(value)) if isinstance(value, Mapping) else {})
for first_class in _COMPILE_TYPED_KEYS:
if first_class in remaining:
compile_config[first_class] = remaining.pop(first_class)
if remaining:
compile_config["extras"] = remaining
elif key == "ltx2_vae_tiling":
pipeline["vae_tiling"] = value
elif key == "config_model_path":
components["config_root"] = value
elif key == "ltx2_refine_enabled":
preset_refine["enabled"] = value
elif key == "ltx2_refine_upsampler_path":
# Empty string means "no upsampler"; keep typed None.
components["upsampler_weights"] = value or None
elif key == "ltx2_refine_lora_path":
# Empty string means "no refine LoRA"; keep typed None.
components["lora_path"] = value or None
elif key == "ltx2_refine_add_noise":
preset_refine["add_noise"] = value
elif key == "ltx2_refine_num_inference_steps":
preset_refine["num_inference_steps"] = value
elif key == "ltx2_refine_guidance_scale":
preset_refine["guidance_scale"] = value
elif key in {"enable_stage_verification", "use_fsdp_inference", "disable_autocast"}:
engine[key] = value
elif key == "override_text_encoder_quant":
@@ -147,6 +184,10 @@ def legacy_from_pretrained_to_config(
if components:
pipeline["components"] = components
if preset_refine:
preset_overrides["refine"] = preset_refine
if preset_overrides:
pipeline["preset_overrides"] = preset_overrides
if experimental:
pipeline["experimental"] = experimental
if pipeline:
@@ -162,12 +203,8 @@ def generator_config_to_fastvideo_args(config: GeneratorConfig | Mapping[str, An
unsupported.append("pipeline.preset")
if normalized.pipeline.preset_version is not None:
unsupported.append("pipeline.preset_version")
if normalized.pipeline.components.config_root is not None:
unsupported.append("pipeline.components.config_root")
if normalized.pipeline.components.vae_weights is not None:
unsupported.append("pipeline.components.vae_weights")
if normalized.pipeline.components.upsampler_weights is not None:
unsupported.append("pipeline.components.upsampler_weights")
if unsupported:
joined = ", ".join(unsupported)
raise NotImplementedError(f"VideoGenerator compatibility adapter does not support {joined} yet")
@@ -191,13 +228,21 @@ def generator_config_to_fastvideo_args(config: GeneratorConfig | Mapping[str, An
"vae_cpu_offload": engine.offload.vae,
"pin_cpu_memory": engine.offload.pin_cpu_memory,
"enable_torch_compile": engine.compile.enabled,
"torch_compile_kwargs": deepcopy(engine.compile.kwargs),
"torch_compile_kwargs": _compile_config_to_torch_kwargs(engine.compile),
"enable_stage_verification": engine.enable_stage_verification,
"use_fsdp_inference": engine.use_fsdp_inference,
"disable_autocast": engine.disable_autocast,
}
if normalized.pipeline.workload_type is not None:
kwargs["workload_type"] = normalized.pipeline.workload_type
if normalized.pipeline.vae_tiling is not None:
kwargs["ltx2_vae_tiling"] = normalized.pipeline.vae_tiling
if engine.compile.text_encoder_enabled is not None:
# ``FastVideoArgs.from_kwargs`` filters to declared fields, so
# this is a no-op on the current legacy path. Emit anyway so the
# realtime runtime (PR 7.6) — which reads from the kwargs dict
# before FastVideoArgs filtering — can pick it up once wired.
kwargs["enable_torch_compile_text_encoder"] = (engine.compile.text_encoder_enabled)
quantization = engine.quantization
if quantization is not None and quantization.text_encoder_quant is not None:
@@ -220,8 +265,18 @@ def generator_config_to_fastvideo_args(config: GeneratorConfig | Mapping[str, An
kwargs["init_weights_from_safetensors"] = components.transformer_weights
if components.transformer_2_weights is not None:
kwargs["init_weights_from_safetensors_2"] = components.transformer_2_weights
if components.config_root is not None:
kwargs["config_model_path"] = components.config_root
if components.upsampler_weights is not None:
kwargs["ltx2_refine_upsampler_path"] = components.upsampler_weights
kwargs.update(deepcopy(normalized.pipeline.preset_overrides))
preset_overrides = deepcopy(normalized.pipeline.preset_overrides)
refine = preset_overrides.pop("refine", None)
if isinstance(refine, Mapping):
for key in _LTX2_REFINE_FLAT_KEYS:
if key in refine:
kwargs[f"ltx2_refine_{key}"] = refine[key]
kwargs.update(preset_overrides)
kwargs.update(deepcopy(normalized.pipeline.experimental))
return FastVideoArgs.from_kwargs(**kwargs)
@@ -271,10 +326,13 @@ def request_to_sampling_param(
) -> SamplingParam:
if request.plan is not None:
raise NotImplementedError("GenerationRequest.plan is not wired into VideoGenerator yet")
if request.state is not None:
raise NotImplementedError("GenerationRequest.state is not wired into VideoGenerator yet")
sampling_param = SamplingParam.from_pretrained(model_path)
if request.state is not None:
_validate_continuation_state(request.state)
sampling_param.continuation_state = request.state
if request.output.return_state:
sampling_param.return_continuation_state = True
updates = explicit_request_updates(request)
for key, value in updates.items():
@@ -316,6 +374,25 @@ def _looks_like_run_or_serve_config(raw: Mapping[str, Any]) -> bool:
return isinstance(raw.get("generator"), Mapping)
def _compile_config_to_torch_kwargs(compile_config: CompileConfig, ) -> dict[str, Any]:
"""Flatten typed ``CompileConfig`` back to a ``torch_compile_kwargs``
dict that the legacy ``FastVideoArgs`` path still expects.
Typed first-class fields (:attr:`backend`, :attr:`fullgraph`,
:attr:`mode`, :attr:`dynamic`) are only emitted when the user set
them explicitly (non-``None``). ``extras`` is merged on top for any
uncommon kwargs.
"""
out: dict[str, Any] = {}
for key in _COMPILE_TYPED_KEYS:
value = getattr(compile_config, key)
if value is not None:
out[key] = value
if compile_config.extras:
out.update(deepcopy(compile_config.extras))
return out
def _sampling_param_to_request_raw(sampling_param: SamplingParam | None, ) -> dict[str, Any]:
if sampling_param is None:
return {}
@@ -476,6 +553,37 @@ def _serialize_generation_request(request: GenerationRequest) -> dict[str, Any]:
_SCHEMA_DEFAULT_UPDATES = _extract_request_updates(config_to_dict(GenerationRequest()))
_KNOWN_CONTINUATION_KINDS: set[str] = set()
def register_continuation_kind(kind: str) -> None:
"""Register a :class:`ContinuationState.kind` as recognized.
PR 7 wires the envelope through; per-kind payload deserializers live
with each model family (e.g. ``fastvideo.pipelines.basic.ltx2.
continuation.LTX2ContinuationState``). The registry lets the
public-API compat layer validate the kind early, before the state
reaches the pipeline.
"""
if not isinstance(kind, str) or not kind:
raise ValueError("ContinuationState kind must be a non-empty string")
_KNOWN_CONTINUATION_KINDS.add(kind)
def _validate_continuation_state(state: ContinuationState) -> None:
if not isinstance(state.kind, str) or not state.kind:
raise ValueError("GenerationRequest.state.kind must be a non-empty string; got "
f"{state.kind!r}")
if not isinstance(state.payload, Mapping):
raise ValueError(f"GenerationRequest.state.payload must be a mapping; got "
f"{type(state.payload).__name__}")
if state.kind not in _KNOWN_CONTINUATION_KINDS:
known = sorted(_KNOWN_CONTINUATION_KINDS)
raise ValueError(f"Unknown ContinuationState kind {state.kind!r}; registered "
f"kinds: {known}. Import the model family that owns this kind "
"(e.g. `import fastvideo.pipelines.basic.ltx2.continuation`) "
"to register it, or drop the state field.")
def _fan_out_batched_input_value(
source_request: GenerationRequest,
@@ -509,6 +617,7 @@ __all__ = [
"load_generator_config_from_file",
"normalize_generation_request",
"normalize_generator_config",
"register_continuation_kind",
"request_to_pipeline_overrides",
"request_to_sampling_param",
]
+4
View File
@@ -15,6 +15,7 @@ class GenerationResult:
samples: Any | None = None
frames: Any | None = None
audio: Any | None = None
audio_sample_rate: int | None = None
size: tuple[int, int, int] | None = None
generation_time: float | None = None
logging_info: Any | None = None
@@ -44,6 +45,7 @@ class GenerationResult:
"samples",
"frames",
"audio",
"audio_sample_rate",
"size",
"generation_time",
"logging_info",
@@ -62,6 +64,7 @@ class GenerationResult:
samples=result.get("samples"),
frames=result.get("frames"),
audio=result.get("audio"),
audio_sample_rate=result.get("audio_sample_rate"),
size=result.get("size"),
generation_time=result.get("generation_time"),
logging_info=result.get("logging_info"),
@@ -80,6 +83,7 @@ class GenerationResult:
"samples": self.samples,
"frames": self.frames,
"audio": self.audio,
"audio_sample_rate": self.audio_sample_rate,
"size": self.size,
"generation_time": self.generation_time,
"logging_info": self.logging_info,
+48 -6
View File
@@ -1,11 +1,16 @@
# SPDX-License-Identifier: Apache-2.0
from __future__ import annotations
import copy
from dataclasses import dataclass, field, fields
from typing import Any
from typing import TYPE_CHECKING, Any
from fastvideo.logger import init_logger
from fastvideo.utils import StoreBoolean
if TYPE_CHECKING:
from fastvideo.api.schema import ContinuationState
logger = init_logger(__name__)
@@ -92,9 +97,13 @@ class SamplingParam:
movement_distance: float | None = None
camera_rotation: str | None = None
# LTX2 multi-modal CFG and STG
ltx2_cfg_scale_video: float = 3.0
ltx2_cfg_scale_audio: float = 7.0
# LTX-2 multi-modal CFG and STG.
# cfg_scale defaults are 1.0 (CFG off) so ``ForwardBatch.__post_init__``
# doesn't force ``do_classifier_free_guidance`` on non-LTX-2 models that
# never override these fields. LTX-2 presets that need text-CFG on set
# them in their ``defaults`` dict (e.g. ``ltx2_base``).
ltx2_cfg_scale_video: float = 1.0
ltx2_cfg_scale_audio: float = 1.0
ltx2_modality_scale_video: float = 3.0
ltx2_modality_scale_audio: float = 3.0
ltx2_rescale_scale: float = 0.7
@@ -103,6 +112,39 @@ class SamplingParam:
ltx2_stg_blocks_video: list[int] = field(default_factory=lambda: [29])
ltx2_stg_blocks_audio: list[int] = field(default_factory=lambda: [29])
# Stable Audio (T2A): clip start/end in seconds. Honored by
# `StableAudioConditioningStage` + `StableAudioDecodingStage`. Other
# families ignore them.
audio_start_in_s: float | None = None
audio_end_in_s: float | None = None
# Stable Audio audio-to-audio (variation):
# `init_audio` -- a path or `[B, C, samples]` waveform at the model
# sample rate; the pipeline encodes it via the VAE
# and uses it as the starting latent.
# `init_audio_strength` -- 0..1, higher = closer to the reference
# (matches the convention of Stability's
# commercial Stable Audio 2.0 UI). 1.0 ~=
# VAE round-trip, 0.0 ~= plain T2A.
# `init_noise_level` -- legacy raw `sigma_max` override (0.3..500,
# higher = more freedom). Kept for callers
# that already use it; prefer `init_audio_strength`.
init_audio: Any = None
init_audio_strength: float | None = None
init_noise_level: float | None = None
# Stable Audio inpainting (RePaint-style): `inpaint_audio` is the
# reference clip, `inpaint_mask` is a [samples] tensor in {0, 1} where
# 1 means *keep the reference* and 0 means *regenerate*.
inpaint_audio: Any = None
inpaint_mask: Any = None
# Continuation state carried across streaming/multi-segment calls.
continuation_state: ContinuationState | None = None
# When True, the pipeline returns a ContinuationState on the result so
# the caller can resume from the generated segment.
return_continuation_state: bool = False
# Misc
save_video: bool = True
return_frames: bool = True
@@ -127,7 +169,7 @@ class SamplingParam:
self.__post_init__()
@classmethod
def from_pretrained(cls, model_path: str) -> "SamplingParam":
def from_pretrained(cls, model_path: str) -> SamplingParam:
sampling_param = cls._from_preset(model_path)
if sampling_param is not None:
return sampling_param
@@ -143,7 +185,7 @@ class SamplingParam:
def _from_preset(
cls,
model_path: str,
) -> "SamplingParam | None":
) -> SamplingParam | None:
"""Build a SamplingParam from preset defaults.
Returns ``None`` when no preset is configured for
+20 -1
View File
@@ -33,8 +33,25 @@ class OffloadConfig:
@dataclass
class CompileConfig:
"""Typed ``torch.compile`` configuration.
``backend``/``fullgraph``/``mode``/``dynamic`` are the four most
common ``torch.compile`` knobs. ``extras`` holds any remaining
``torch.compile`` kwargs (e.g. ``options``, ``disable``).
"""
enabled: bool = False
kwargs: dict[str, Any] = field(default_factory=dict)
text_encoder_enabled: bool | None = None
"""Whether ``torch.compile`` is applied to the text encoder. ``None``
keeps the runtime default. The public ``FastVideoArgs`` adapter does
not yet consume this flag; reserved so the realtime runtime upstream
(PR 7.6) has a typed home for its ``enable_torch_compile_text_encoder``
kwarg without routing through ``pipeline.experimental``."""
backend: str | None = None
fullgraph: bool | None = None
mode: str | None = None
dynamic: bool | None = None
extras: dict[str, Any] = field(default_factory=dict)
@dataclass
@@ -76,6 +93,8 @@ class PipelineSelection:
preset: str | None = None
preset_version: int | None = None
components: ComponentConfig = field(default_factory=ComponentConfig)
vae_tiling: bool | None = None
"""Tile-based VAE decode. ``None`` keeps the model's default."""
preset_overrides: dict[str, Any] = field(default_factory=dict)
experimental: dict[str, Any] = field(default_factory=dict)
+58
View File
@@ -0,0 +1,58 @@
# `fastvideo/attention/` — Attention Backends
**Generated:** 2026-05-02
Backend registry + selector wrapping FlashAttn / SageAttn / SageAttn3 / SDPA / VSA / VMoBA / SLA / BSA.
## Layout
```
attention/
├── __init__.py # Exports DistributedAttention, LocalAttention, get_attn_backend
├── layer.py # DistributedAttention, DistributedAttention_VSA, LocalAttention
├── selector.py # get_attn_backend (cached) + env-var override
├── backends/
│ ├── abstract.py # AttentionBackend / AttentionMetadata / AttentionMetadataBuilder
│ ├── flash_attn.py # FA2/FA3
│ ├── sage_attn.py # SageAttention v1
│ ├── sage_attn3.py # SageAttention v3
│ ├── sdpa.py # torch SDPA fallback
│ ├── video_sparse_attn.py # VSA (paper: Video Sparse Attention)
│ ├── vmoba.py # Video-MoBA
│ ├── sla.py # Sliding-window (STA)
│ └── bsa_attn.py # Block-sparse
└── utils/
├── flash_attn_cute.py
└── flash_attn_no_pad.py
```
## Selection Order
`get_attn_backend()` resolves via:
1. Env-var override `FASTVIDEO_ATTENTION_BACKEND` (see `STR_BACKEND_ENV_VAR` in `fastvideo/utils.py`).
2. Per-platform default from `fastvideo/platforms/`.
3. Heuristic fallback to SDPA.
The result is `@lru_cache`d. Tests that need a specific backend must use the
`global_force_attn_backend(...)` context manager from `selector.py`, never set
the env var mid-process.
## Adding a Backend
1. Subclass `AttentionBackend` in `backends/<name>.py`.
2. Implement `AttentionMetadata` + `AttentionMetadataBuilder` for the new path.
3. Register the enum value in `fastvideo/platforms/interface.py` (`AttentionBackendEnum`).
4. Wire string → class resolution in `selector.py`.
5. Verify the new backend works with `DistributedAttention` (sequence parallel)
and `LocalAttention` (single-rank). If it cannot support SP, document the
gap in the backend file's module docstring.
## Anti-Patterns
- Calling `torch.nn.functional.scaled_dot_product_attention` directly inside a
model's forward — go through `DistributedAttention` / `LocalAttention`.
- Reading `os.environ[STR_BACKEND_ENV_VAR]` from arbitrary call sites. Use
`get_env_variable_attn_backend()`.
- Caching backend instances per-module. The selector cache is process-wide; do
not duplicate it.
+1 -1
View File
@@ -2,7 +2,6 @@
import torch
import torch.nn.functional as F
from flash_attn import flash_attn_func as flash_attn_2_func
from dataclasses import dataclass
try:
@@ -18,6 +17,7 @@ except ImportError:
flash_attn_func = flash_attn_3_func
fa_version = "3"
except ImportError:
from flash_attn import flash_attn_func as flash_attn_2_func
flash_attn_func = flash_attn_2_func
fa_version = "2"
+1 -1
View File
@@ -405,7 +405,7 @@ class SageSLAAttentionImpl(AttentionImpl, nn.Module):
if not SAGESLA_ENABLED:
raise ImportError("SageSLA requires spas_sage_attn. "
"Install with: pip install git+https://github.com/thu-ml/SpargeAttn.git")
"Install with: uv pip install git+https://github.com/thu-ml/SpargeAttn.git")
assert head_size in [64, 128], f"SageSLA requires head_size in [64, 128], got {head_size}"
+53
View File
@@ -0,0 +1,53 @@
# `fastvideo/configs/` — Config-Driven Model Registry
**Generated:** 2026-05-02
Two layers of dataclass configs feed every pipeline: **arch configs** (what the model is) and **pipeline configs** (how to run it).
## Layout
```
configs/
├── configs.py # Dataset / loader enums (DatasetType, VideoLoaderType)
├── utils.py # update_config_from_args, shallow_asdict helpers
├── backend/ # Attention backend defaults
├── models/
│ ├── base.py # ModelConfig ABC
│ ├── dits/ # DiTConfig per model (wanvideo, ltx2, hunyuan, ...)
│ ├── vaes/ # VAEConfig per model
│ ├── encoders/ # EncoderConfig (t5, clip, llama, qwen2_5, gemma, siglip, ...)
│ ├── upsamplers/ # UpsamplerConfig (hunyuan15)
│ └── audio/ # Audio-model configs (ltx2_audio_vae, ...)
├── pipelines/
│ ├── base.py # PipelineConfig ABC + (de)serialization
│ └── <model>.py # Concrete configs (HunyuanConfig, WanT2V480PConfig, ...)
└── *.json # Frozen reference configs for shipped models
```
## How Configs Hook Into the Registry
`fastvideo/registry.py` imports every concrete `PipelineConfig` and exposes
`get_pipeline_config_cls_from_name(...)`. Adding a new pipeline config requires:
1. Subclass `PipelineConfig` in `pipelines/<model>.py`.
2. Reference its component arch configs (DiT / VAE / encoder / upsampler).
3. Add the import + name mapping in `fastvideo/registry.py`.
Configs that do not appear in `registry.py` are unreachable from `VideoGenerator`.
## Arch vs Pipeline — Where Does This Field Go?
| Field type | Lives on |
|-----------|----------|
| Architecture constants (hidden dim, num heads, layer count) | `configs/models/<role>/<model>.py` |
| Default sampling params (steps, cfg, shift, fps) | `configs/pipelines/<model>.py` |
| Runtime overrides (precision, sp_size, tp_size, attention backend) | `configs/pipelines/base.py` defaults + CLI flags via `fastvideo_args.py` |
| `param_names_mapping` for HF → FastVideo state-dict | Arch config (lives with the model definition) |
If a knob is tunable per inference call → `SamplingParam`, not `PipelineConfig`.
## Anti-Patterns
- Hard-coding architecture constants inside model classes — always read from the arch config.
- Using `argparse` directly here. Configs deserialize from dicts via `update_config_from_args`.
- Importing from `fastvideo.pipelines` here. Configs are the lower layer; the dependency is one-way.
-2
View File
@@ -48,8 +48,6 @@ class ModelConfig:
for key, value in source_model_dict.items():
if key in valid_fields:
setattr(arch_config, key, value)
else:
raise AttributeError(f"{type(arch_config).__name__} has no field '{key}'")
if hasattr(arch_config, "__post_init__"):
arch_config.__post_init__()
+4 -1
View File
@@ -5,11 +5,14 @@ from fastvideo.configs.models.dits.hunyuanvideo import HunyuanVideoConfig
from fastvideo.configs.models.dits.hunyuanvideo15 import HunyuanVideo15Config
from fastvideo.configs.models.dits.longcat import LongCatVideoConfig
from fastvideo.configs.models.dits.ltx2 import LTX2VideoConfig
from fastvideo.configs.models.dits.magi_human import MagiHumanVideoConfig
from fastvideo.configs.models.dits.stable_audio import StableAudioConfig
from fastvideo.configs.models.dits.wanvideo import WanVideoConfig
from fastvideo.configs.models.dits.hyworld import HYWorldConfig
from fastvideo.configs.models.dits.kandinsky5 import Kandinsky5VideoConfig
__all__ = [
"HunyuanVideoConfig", "HunyuanVideo15Config", "HunyuanGameCraftConfig", "WanVideoConfig", "CosmosVideoConfig",
"Cosmos25VideoConfig", "LongCatVideoConfig", "LTX2VideoConfig", "HYWorldConfig", "Kandinsky5VideoConfig"
"Cosmos25VideoConfig", "LongCatVideoConfig", "LTX2VideoConfig", "HYWorldConfig", "Kandinsky5VideoConfig",
"MagiHumanVideoConfig", "StableAudioConfig"
]
+2 -1
View File
@@ -50,7 +50,8 @@ class CosmosArchConfig(DiTArchConfig):
})
# Cosmos-specific config parameters based on transformer_cosmos.py
in_channels: int = 16
# in_channels includes the condition_mask channel (16 latent + 1 cond = 17)
in_channels: int = 17
out_channels: int = 16
num_attention_heads: int = 16
attention_head_dim: int = 128
+110
View File
@@ -0,0 +1,110 @@
# SPDX-License-Identifier: Apache-2.0
"""Architecture / model config for the daVinci-MagiHuman DiT.
The MagiHuman base DiT is a 15B-parameter single-stream transformer that
jointly denoises video, audio, and text tokens in one flat sequence. Layout
details verified against GAIR/daVinci-MagiHuman's base/ shards (2026-04-24).
This file captures only configuration. The module implementation lives in
fastvideo/models/dits/magi_human.py and the pipeline wiring in
fastvideo/pipelines/basic/magi_human/.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from fastvideo.configs.models.dits.base import DiTArchConfig, DiTConfig
def _is_block_layer(n: str, m) -> bool:
# Match "block.layers.<idx>" — the FSDP shard boundary for MagiHuman.
parts = n.split(".")
return (len(parts) >= 3 and parts[0] == "block" and parts[1] == "layers" and str.isdigit(parts[2]))
@dataclass
class MagiHumanArchConfig(DiTArchConfig):
"""MagiHuman base DiT architecture constants.
**Scope contract:** fields here must match the `transformer/config.json`
emitted by `scripts/checkpoint_conversion/convert_magi_human_to_diffusers.py`
1:1, and both are sourced from the upstream Python reference
`inference/common/config.py::ModelConfig` (the HF root `config.json`
is empty so the Python source is canonical). Pipeline-level knobs
(VAE stride, fps, num_inference_steps, CFG scales, flow_shift,
t5_gemma_target_length) and data-proxy knobs (coords_style,
frame_receptive_field, ref_audio_offset, text_offset) live on
`MagiHumanBaseConfig`, NOT here.
`param_names_mapping` is intentionally empty: the FastVideo implementation
keeps the same module tree as the reference (`adapter.*`,
`block.layers.<i>.*`, `final_linear_{video,audio}.*`,
`final_norm_{video,audio}.*`), so converted weights load directly.
"""
_fsdp_shard_conditions: list = field(default_factory=lambda: [_is_block_layer])
# No renames needed — the FastVideo module mirrors the reference names.
param_names_mapping: dict = field(default_factory=dict)
reverse_param_names_mapping: dict = field(default_factory=dict)
lora_param_names_mapping: dict = field(default_factory=dict)
# --- transformer shape ---
num_layers: int = 40
hidden_size: int = 5120
head_dim: int = 128
num_query_groups: int = 8 # num_heads_kv (GQA)
# --- modality channels ---
# video_in_channels = z_dim (48) * patch_size product (1*2*2=4), so the
# embedder receives 192 per token. text_in_channels is T5Gemma-9B's
# encoder hidden size.
video_in_channels: int = 192
audio_in_channels: int = 64
text_in_channels: int = 3584
# --- block-level architecture switches ---
# Sandwich MoE: first and last 4 layers have per-modality experts
# (video/audio/text), middle layers share a single set of weights.
mm_layers: tuple[int, ...] = (0, 1, 2, 3, 36, 37, 38, 39)
local_attn_layers: tuple[int, ...] = ()
gelu7_layers: tuple[int, ...] = (0, 1, 2, 3)
post_norm_layers: tuple[int, ...] = ()
enable_attn_gating: bool = True
activation_type: str = "swiglu7"
# --- DiT patching (upstream `ModelConfig`-equivalent; NOT the VAE
# stride, which is pipeline-level). ---
patch_size: tuple[int, int, int] = (1, 2, 2)
spatial_rope_interpolation: str = "extra"
# --- TReAD (token routing + early drop). Flattened from the upstream
# nested `tread_config` dict so it round-trips through
# `update_model_arch` cleanly. ---
tread_selection_rate: float = 0.5
tread_start_layer_idx: int = 2
tread_end_layer_idx: int = 25
# --- derived fields (populated in __post_init__) ---
num_attention_heads: int = 0 # hidden_size / head_dim
num_heads_kv: int = 0 # == num_query_groups
in_channels: int = 0 # mirror of video_in_channels (FastVideo contract)
out_channels: int = 0 # mirror of video_in_channels
def __post_init__(self) -> None:
super().__post_init__()
self.num_attention_heads = self.hidden_size // self.head_dim
self.num_heads_kv = self.num_query_groups
self.in_channels = self.video_in_channels
self.out_channels = self.video_in_channels
# num_channels_latents is the VAE latent z_dim (48 for Wan 2.2 TI2V-5B).
# We don't declare z_dim on the arch config (it's a VAE property),
# but we still set num_channels_latents for the BaseDiT contract.
self.num_channels_latents = 48
@dataclass
class MagiHumanVideoConfig(DiTConfig):
arch_config: DiTArchConfig = field(default_factory=MagiHumanArchConfig)
prefix: str = "magi_human"
@@ -0,0 +1,76 @@
# SPDX-License-Identifier: Apache-2.0
"""Config for the Stable Audio Open 1.0 DiT.
Note: the SA pipeline bypasses the standard `ComposedPipelineBase`
component loader because the published HF repo ships a single monolithic
`model.safetensors` (no Diffusers-style `model_index.json` or
per-subfolder layout). The arch fields and `param_names_mapping` here
document the architecture and key remap so the same conventions used by
the rest of the DiT family apply (FSDP shard conditions, supported
attention backends, future loader integrations) — they are not currently
consumed by `fastvideo/models/loader/fsdp_load.py` for SA.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from fastvideo.configs.models.dits.base import DiTArchConfig, DiTConfig
from fastvideo.platforms import AttentionBackendEnum
def _is_transformer_layer(n: str, m) -> bool:
# Matches `transformer.layers.{i}` in the SA DiT module tree.
parts = n.split(".")
return (len(parts) >= 3 and parts[-3] == "transformer" and parts[-2] == "layers" and parts[-1].isdigit())
@dataclass
class StableAudioArchConfig(DiTArchConfig):
_fsdp_shard_conditions: list = field(default_factory=lambda: [_is_transformer_layer])
# SA's checkpoint is `stable_audio_tools` raw format (not Diffusers),
# so the only remaps are: strip the `model.model.` host-pipeline
# prefix, and rename `nn.LayerNorm`'s `gamma`/`beta` to torch's
# canonical `weight`/`bias`. Linear / cross-attention naming already
# matches FastVideo's conventions, so no further remap is needed.
param_names_mapping: dict = field(
default_factory=lambda: {
r"^model\.model\.(.*?)\.gamma$": r"\1.weight",
r"^model\.model\.(.*?)\.beta$": r"\1.bias",
r"^model\.model\.(.*)$": r"\1",
})
# SA only supports backends compatible with single-GPU LocalAttention.
_supported_attention_backends: tuple[AttentionBackendEnum, ...] = (
AttentionBackendEnum.FLASH_ATTN,
AttentionBackendEnum.TORCH_SDPA,
)
# Architecture constants (from the published `model_config.json` for
# `stabilityai/stable-audio-open-1.0`).
io_channels: int = 64
embed_dim: int = 1536
depth: int = 24
num_attention_heads: int = 24
cond_token_dim: int = 768
global_cond_dim: int = 1536
project_cond_tokens: bool = False
project_global_cond: bool = True
# Set to "ln" to wrap attention Q/K in LayerNorm (used by
# `stable-audio-open-small`; absent in the 1.0 base).
qk_norm: str | None = None
def __post_init__(self) -> None:
super().__post_init__()
self.hidden_size = self.embed_dim
self.in_channels = self.io_channels
self.out_channels = self.io_channels
self.num_channels_latents = self.io_channels
self.attention_head_dim = self.embed_dim // self.num_attention_heads
@dataclass
class StableAudioConfig(DiTConfig):
arch_config: DiTArchConfig = field(default_factory=StableAudioArchConfig)
prefix: str = "StableAudio"
@@ -7,9 +7,13 @@ from fastvideo.configs.models.encoders.qwen2_5 import Qwen2_5_VLConfig
from fastvideo.configs.models.encoders.siglip import SiglipVisionConfig
from fastvideo.configs.models.encoders.reason1 import Reason1ArchConfig, Reason1Config
from fastvideo.configs.models.encoders.gemma import LTX2GemmaConfig
from fastvideo.configs.models.encoders.stable_audio_conditioner import (StableAudioConditionerArchConfig,
StableAudioConditionerConfig)
from fastvideo.configs.models.encoders.t5gemma import T5GemmaEncoderConfig
__all__ = [
"EncoderConfig", "TextEncoderConfig", "ImageEncoderConfig", "BaseEncoderOutput", "CLIPTextConfig",
"CLIPVisionConfig", "WAN2_1ControlCLIPVisionConfig", "LlamaConfig", "T5Config", "T5LargeConfig", "Qwen2_5_VLConfig",
"Reason1ArchConfig", "Reason1Config", "LTX2GemmaConfig", "SiglipVisionConfig"
"Reason1ArchConfig", "Reason1Config", "LTX2GemmaConfig", "SiglipVisionConfig", "StableAudioConditionerArchConfig",
"StableAudioConditionerConfig", "T5GemmaEncoderConfig"
]
@@ -0,0 +1,80 @@
# SPDX-License-Identifier: Apache-2.0
"""Config for the Stable Audio Open 1.0 multi-conditioner.
The conditioner bundles three sub-conditioners — a T5 text encoder
(prompt) and two NumberConditioners (`seconds_start` / `seconds_total`)
— into the (cross_attn_cond, cross_attn_mask, global_embed) triple the
DiT consumes. The architecture is fully specified by the official
`stable_audio_tools` `MultiConditioner` config; the constants here
mirror that.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from fastvideo.configs.models.base import ArchConfig
from fastvideo.configs.models.encoders.base import (EncoderArchConfig, EncoderConfig)
def _default_configs() -> list[dict]:
"""Default = `stable-audio-open-1.0`'s three sub-conditioners."""
return [
{
"id": "prompt",
"type": "t5",
"config": {
"t5_model_name": "t5-base",
"max_length": 128
}
},
{
"id": "seconds_start",
"type": "number",
"config": {
"min_val": 0,
"max_val": 512
}
},
{
"id": "seconds_total",
"type": "number",
"config": {
"min_val": 0,
"max_val": 512
}
},
]
@dataclass
class StableAudioConditionerArchConfig(EncoderArchConfig):
architectures: list[str] = field(default_factory=lambda: ["StableAudioMultiConditioner"])
# Shared embedding width across all sub-conditioners (T5 last-hidden
# dim and NumberEmbedder feature dim both = `cond_dim`).
cond_dim: int = 768
# Sub-conditioner identifiers. Order in `cross_attention_cond_ids`
# is the concat order for the cross-attn token sequence; order in
# `global_cond_ids` is the concat order for the global FiLM-style
# embedding.
cross_attention_cond_ids: tuple[str, ...] = ("prompt", "seconds_start", "seconds_total")
global_cond_ids: tuple[str, ...] = ("seconds_start", "seconds_total")
# Per-sub-conditioner spec list (mirrors upstream
# `model_config.json.model.conditioning.configs`). Each entry is
# `{"id": ..., "type": "t5"|"number", "config": {...}}`. The default
# matches `stable-audio-open-1.0`; SA-small overrides via the
# `conditioner/config.json` shipped in the converted repo.
configs: list = field(default_factory=_default_configs)
# Match official `stable_audio_tools/models/conditioners.py:334`:
# T5 is loaded directly in fp16.
t5_dtype: str = "float16"
@dataclass
class StableAudioConditionerConfig(EncoderConfig):
arch_config: ArchConfig = field(default_factory=StableAudioConditionerArchConfig)
prefix: str = "stable_audio_conditioner"
+8
View File
@@ -41,6 +41,14 @@ class T5ArchConfig(TextEncoderArchConfig):
text_len: int = 512
dtype: str | None = None
gradient_checkpointing: bool = False
# Extra fields present in upstream HF T5Config but unused by FastVideo's
# encoder. Declared here so `update_model_arch` doesn't reject them when
# loading repos like `stabilityai/stable-audio-open-1.0` that ship the
# full HF config.
n_positions: int = 512
decoder_start_token_id: int = 0
output_past: bool = True
task_specific_params: dict | None = None
stacked_params_mapping: list[tuple[str, str, str]] = field(default_factory=lambda: [
# (param_name, shard_name, shard_id)
(".qkv_proj", ".q", "q"),
@@ -0,0 +1,71 @@
# SPDX-License-Identifier: Apache-2.0
"""Config for the T5-Gemma encoder used by daVinci-MagiHuman.
The reference pipeline uses `transformers.models.t5gemma.T5GemmaEncoderModel`
on `google/t5gemma-9b-9b-ul2`. That is a gated Google repository, so the
encoder weights are not bundled inside GAIR/daVinci-MagiHuman; they are
loaded from the T5-Gemma HF repo directly.
Encoder shape (verified from google/t5gemma-9b-9b-ul2/config.json):
layers=42, hidden=3584, heads=16, kv_heads=8, head_dim=256,
intermediate=14336, rope_theta=10000.0, max_pos=8192,
layer_types alternate sliding_attention / full_attention.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from fastvideo.configs.models.encoders.base import (
TextEncoderArchConfig,
TextEncoderConfig,
)
def _is_t5gemma_model(n: str, m) -> bool:
return n.endswith("t5gemma_model") or n.endswith("_t5gemma_model")
@dataclass
class T5GemmaEncoderArchConfig(TextEncoderArchConfig):
architectures: list[str] = field(default_factory=lambda: ["T5GemmaEncoderModel"])
hidden_size: int = 3584
num_hidden_layers: int = 42
num_attention_heads: int = 16
num_key_value_heads: int = 8
head_dim: int = 256
intermediate_size: int = 14336
max_position_embeddings: int = 8192
rope_theta: float = 10000.0
vocab_size: int = 256000
# MagiHuman fixes prompt embed length at 640 via pad_or_trim.
text_len: int = 640
pad_token_id: int = 0
eos_token_id: int = 1
# Path to the upstream gated repo. When set, the FastVideo loader will
# pull the encoder directly via `T5GemmaEncoderModel.from_pretrained`.
t5gemma_model_path: str = "google/t5gemma-9b-9b-ul2"
t5gemma_dtype: str = "bfloat16"
_fsdp_shard_conditions: list = field(default_factory=lambda: [_is_t5gemma_model])
def __post_init__(self) -> None:
super().__post_init__()
# WHY: upstream `t5_gemma_model.py:25` tokenizes without
# padding/max_length, then `prompt_process.py` pad_or_trim-s the
# encoded states. Keep only tensor return here so
# MagiHumanLatentPreparationStage can pad/trim post-encode while
# preserving the real original prompt length.
self.tokenizer_kwargs.pop("truncation", None)
self.tokenizer_kwargs.pop("max_length", None)
self.tokenizer_kwargs.pop("padding", None)
@dataclass
class T5GemmaEncoderConfig(TextEncoderConfig):
arch_config: TextEncoderArchConfig = field(default_factory=T5GemmaEncoderArchConfig)
prefix: str = "t5gemma"
@@ -5,6 +5,7 @@ from fastvideo.configs.models.vaes.gen3cvae import Gen3CVAEConfig
from fastvideo.configs.models.vaes.hunyuanvae import HunyuanVAEConfig
from fastvideo.configs.models.vaes.hunyuan15vae import Hunyuan15VAEConfig
from fastvideo.configs.models.vaes.ltx2vae import LTX2VAEConfig
from fastvideo.configs.models.vaes.oobleck import OobleckVAEArchConfig, OobleckVAEConfig
from fastvideo.configs.models.vaes.wanvae import WanVAEConfig
__all__ = [
@@ -16,4 +17,6 @@ __all__ = [
"Gen3CVAEConfig",
"Hunyuan15VAEConfig",
"LTX2VAEConfig",
"OobleckVAEArchConfig",
"OobleckVAEConfig",
]
+68
View File
@@ -0,0 +1,68 @@
# SPDX-License-Identifier: Apache-2.0
"""Config for the Stable Audio Open 1.0 "Oobleck" VAE.
Mirrors the per-channel `vae/config.json` shipped in
`stabilityai/stable-audio-open-1.0` 1:1 (see
`fastvideo/models/vaes/oobleck.py::OobleckVAE.from_pretrained`, which
constructs the VAE from these fields). Inherits the FastVideo VAEConfig
base so the standard `load_encoder` / `load_decoder` flags + tiling
knobs apply.
Naming: the VAE architecture is officially "Oobleck" (per Stability
AI's stable-audio-tools) — the surrounding model family is "Stable
Audio Open 1.0". This config is named after the architecture
(`OobleckVAEConfig`) since the same VAE is shared across Stable Audio
checkpoints; downstream pipelines reference it by its arch name, not
by a host-pipeline name.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from fastvideo.configs.models.vaes.base import VAEArchConfig, VAEConfig
@dataclass
class OobleckVAEArchConfig(VAEArchConfig):
"""Stable Audio Open 1.0 VAE architecture constants."""
architectures: list[str] = field(default_factory=lambda: ["AutoencoderOobleck"])
# From stabilityai/stable-audio-open-1.0/vae/config.json.
encoder_hidden_size: int = 128
downsampling_ratios: list[int] = field(default_factory=lambda: [2, 4, 4, 8, 8])
channel_multiples: list[int] = field(default_factory=lambda: [1, 2, 4, 8, 16])
decoder_channels: int = 128
decoder_input_channels: int = 64
audio_channels: int = 2 # stereo
sampling_rate: int = 44100
@dataclass
class OobleckVAEConfig(VAEConfig):
"""FastVideo VAE config wrapping the Oobleck arch.
Audio VAEs don't use the temporal/spatial tiling defaults that the
base VAEConfig is shaped for (those exist for video VAEs); they are
retained but irrelevant for audio.
"""
arch_config: VAEArchConfig = field(default_factory=OobleckVAEArchConfig)
# Audio is 1-D, so the video-VAE tiling defaults are inert. Disable
# them so callers don't accidentally trip on tile-stride math built
# for spatial tensors.
use_tiling: bool = False
use_temporal_tiling: bool = False
use_parallel_tiling: bool = False
# Where the FastVideo loader / pipeline-glue wrapper should fetch
# weights from when no local path is supplied. Gated repo — caller's
# HF token must have accepted terms on
# https://huggingface.co/stabilityai/stable-audio-open-1.0.
pretrained_path: str = "stabilityai/stable-audio-open-1.0"
pretrained_subfolder: str = "vae"
# Match official `stable_audio_tools`: VAE runs in fp16 (the
# `pretransform.model_half` path in
# `stable_audio_tools/models/pretransforms.py`).
pretrained_dtype: str = "float16"
+1 -1
View File
@@ -5,7 +5,7 @@ from fastvideo.configs.pipelines.hunyuan import FastHunyuanConfig, HunyuanConfig
from fastvideo.configs.pipelines.hunyuan15 import Hunyuan15T2V480PConfig, Hunyuan15T2V720PConfig
from fastvideo.configs.pipelines.hunyuangamecraft import HunyuanGameCraftPipelineConfig
from fastvideo.configs.pipelines.hyworld import HYWorldConfig
from fastvideo.configs.pipelines.ltx2 import LTX2T2VConfig
from fastvideo.pipelines.basic.ltx2.pipeline_configs import LTX2T2VConfig
from fastvideo.registry import get_pipeline_config_cls_from_name
from fastvideo.configs.pipelines.wan import (SelfForcingWanT2V480PConfig, WanI2V480PConfig, WanI2V720PConfig,
WanT2V480PConfig, WanT2V720PConfig)
+27 -2
View File
@@ -11,10 +11,34 @@ import torch
from fastvideo.configs.models import DiTConfig, VAEConfig
from fastvideo.configs.models.dits.base import DiTArchConfig
from fastvideo.configs.models.encoders import BaseEncoderOutput, T5Config
from fastvideo.configs.models.encoders.base import TextEncoderArchConfig
from fastvideo.configs.models.encoders.t5 import T5ArchConfig
from fastvideo.configs.models.vaes import WanVAEConfig
from fastvideo.configs.pipelines.base import PipelineConfig
@dataclass
class LongCatT5ArchConfig(T5ArchConfig):
"""T5 arch that pads tokenizer output to ``max_length``.
LongCat's denoising stage concatenates positive and negative
attention masks along the batch dimension for CFG, which requires
uniform seq length. The shared :class:`T5ArchConfig` dropped the
``"padding": "max_length"`` tokenizer kwarg so other DiTs could run
with variable-length masks; LongCat still needs the uniform
contract.
"""
def __post_init__(self) -> None:
super().__post_init__()
self.tokenizer_kwargs["padding"] = "max_length"
@dataclass
class LongCatT5Config(T5Config):
arch_config: TextEncoderArchConfig = field(default_factory=LongCatT5ArchConfig)
@dataclass
class LongCatDiTArchConfig(DiTArchConfig):
"""Extended DiTArchConfig with LongCat-specific fields."""
@@ -103,8 +127,9 @@ class LongCatT2V480PConfig(PipelineConfig):
vae_precision: str = "bf16"
text_encoder_precisions: tuple[str, ...] = field(default_factory=lambda: ("bf16", ))
# Text encoding (UMT5 uses T5-like config; postprocess to fixed 512)
text_encoder_configs: tuple[T5Config, ...] = field(default_factory=lambda: (T5Config(), ))
# UMT5 uses T5-like config; postprocess pads to 512. LongCatT5Config
# restores ``padding="max_length"`` for the CFG concat contract.
text_encoder_configs: tuple[T5Config, ...] = field(default_factory=lambda: (LongCatT5Config(), ))
preprocess_text_funcs: tuple[Callable[[str], str], ...] = field(default_factory=lambda: (longcat_preprocess_text, ))
postprocess_text_funcs: tuple[Callable[[BaseEncoderOutput], torch.Tensor],
...] = field(default_factory=lambda: (umt5_postprocess_text, ))
@@ -0,0 +1,67 @@
# SPDX-License-Identifier: Apache-2.0
"""`PipelineConfig` for Stable Audio Open 1.0."""
from __future__ import annotations
from dataclasses import dataclass, field
from fastvideo.configs.models import DiTConfig, VAEConfig
from fastvideo.configs.models.dits import StableAudioConfig
from fastvideo.configs.models.vaes import OobleckVAEConfig
from fastvideo.configs.pipelines.base import PipelineConfig
@dataclass
class StableAudioT2AConfig(PipelineConfig):
"""Stable Audio Open 1.0 pipeline config."""
dit_config: DiTConfig = field(default_factory=StableAudioConfig)
# Standard `TransformerLoader` reads `dit_precision`; default in
# `PipelineConfig` is bf16, but we want fp16 to match official.
dit_precision: str = "fp16"
vae_config: VAEConfig = field(default_factory=OobleckVAEConfig)
vae_tiling: bool = False
vae_sp: bool = False
# `StableAudioMultiConditioner` owns its own T5; zero out the
# parent's text-encoder slots so the length-equality validator passes.
text_encoder_configs: tuple = field(default_factory=tuple)
preprocess_text_funcs: tuple = field(default_factory=tuple)
postprocess_text_funcs: tuple = field(default_factory=tuple)
num_inference_steps: int = 100
guidance_scale: float = 7.0
audio_end_in_s: float = 10.0 # short-clip default
audio_start_in_s: float = 0.0
sampling_rate: int = 44100
audio_channels: int = 2
# Stable Audio Open 1.0 was trained at a fixed 2,097,152-sample
# window (= 2097152 / 44100 ≈ 47.55s). Anything past this is
# silently truncated by the post-decode slice — validate up-front.
sample_size: int = 2097152
max_audio_duration_s: float = 2097152 / 44100
# Match the official `stable_audio_tools` defaults (`model_half=True`
# in `run_gradio.py`), which loads the DiT, VAE, and T5 in fp16 and
# wraps T5 forward in `autocast(fp16)`. fp16 is also a hard
# requirement for FlashAttention-2 / FA-3.
precision: str = "fp16"
vae_precision: str = "fp16"
text_encoder_precisions: tuple[str, ...] = field(default_factory=tuple)
def __post_init__(self) -> None:
# A2A needs encode; load both halves for either path.
self.vae_config.load_encoder = True
self.vae_config.load_decoder = True
@dataclass
class StableAudioOpenSmallConfig(StableAudioT2AConfig):
"""`stable-audio-open-small` overrides: shorter training window
(524288 samples ≈ 11.89s @ 44.1 kHz) and a faster default sampler
config carried by the small preset.
"""
sample_size: int = 524288
max_audio_duration_s: float = 524288 / 44100
audio_end_in_s: float = 6.0 # short-clip default suitable for the small window
+3
View File
@@ -3,6 +3,8 @@
from fastvideo.entrypoints.cli.cli_types import CLISubcommand
from fastvideo.entrypoints.cli.generate import cmd_init as generate_cmd_init
from fastvideo.utils import FlexibleArgumentParser
from fastvideo.entrypoints.cli.router_serve import (
cmd_init as router_serve_cmd_init, )
from fastvideo.entrypoints.cli.serve import cmd_init as serve_cmd_init
from fastvideo.entrypoints.cli.bench import cmd_init as bench_cmd_init
@@ -12,6 +14,7 @@ def cmd_init() -> list[CLISubcommand]:
commands = []
commands.extend(generate_cmd_init())
commands.extend(serve_cmd_init())
commands.extend(router_serve_cmd_init())
commands.extend(bench_cmd_init())
return commands
+115
View File
@@ -0,0 +1,115 @@
# SPDX-License-Identifier: Apache-2.0
"""``fastvideo router-serve`` CLI subcommand.
Launches the streaming router from a YAML config. Separate from
``fastvideo serve`` because the router is an orthogonal process: it
fronts one or more running servers rather than hosting a generator
itself.
"""
from __future__ import annotations
import argparse
import os
from typing import cast
from fastvideo.api.parser import load_raw_config
from fastvideo.entrypoints.cli.cli_types import CLISubcommand
from fastvideo.entrypoints.streaming.router.config import (
ReplicaEndpoint,
RouterConfig,
)
from fastvideo.logger import init_logger
from fastvideo.utils import FlexibleArgumentParser
logger = init_logger(__name__)
class RouterServeSubcommand(CLISubcommand):
"""Start the multi-replica WebSocket router."""
def __init__(self) -> None:
self.name = "router-serve"
super().__init__()
def cmd(self, args: argparse.Namespace) -> None:
config = _load_router_config(args.config)
logger.info(
"router listening on %s:%d (%d replicas, %d primary)",
config.host,
config.port,
len(config.replicas),
sum(1 for r in config.replicas if r.primary),
)
from fastvideo.entrypoints.streaming.router.main import run_router
run_router(config)
def validate(self, args: argparse.Namespace) -> None:
if not args.config:
raise ValueError("fastvideo router-serve requires --config PATH")
if not os.path.exists(args.config):
raise ValueError(f"Router config file not found: {args.config}")
def subparser_init(
self,
subparsers: argparse._SubParsersAction,
) -> FlexibleArgumentParser:
parser = subparsers.add_parser(
"router-serve",
help="Start the streaming router (multi-replica load balancer)",
usage="fastvideo router-serve --config ROUTER_CONFIG",
)
parser.add_argument(
"--config",
type=str,
default="",
required=False,
help="Path to a YAML/JSON router config. Required.",
)
return cast(FlexibleArgumentParser, parser)
def _load_router_config(path: str) -> RouterConfig:
raw = load_raw_config(path)
router_raw = raw.get("router") if isinstance(raw, dict) else None
if not isinstance(router_raw, dict):
raise ValueError(f"Router config {path!r} must have a top-level `router:` block")
replicas_raw = router_raw.get("replicas", [])
if not isinstance(replicas_raw, list):
raise ValueError(f"router.replicas must be a list, got {type(replicas_raw).__name__}")
replicas = []
for i, r in enumerate(replicas_raw):
if not isinstance(r, dict):
raise ValueError(f"router.replicas[{i}] must be a mapping, got {type(r).__name__}")
url = r.get("url")
if not url:
raise ValueError(f"router.replicas[{i}] is missing required key 'url'")
replicas.append(
ReplicaEndpoint(
url=url,
name=r.get("name"),
primary=bool(r.get("primary", False)),
weight=float(r.get("weight", 1.0)),
))
if not replicas:
raise ValueError("Router config must list at least one replica under `router.replicas`")
health_check = router_raw.get("health_check") or {}
return RouterConfig(
host=str(router_raw.get("host", "0.0.0.0")),
port=int(router_raw.get("port", 9000)),
replicas=replicas,
health_check_path=str(health_check.get("path", "/health")),
health_check_interval_seconds=float(health_check.get("interval_seconds", 5.0)),
health_check_timeout_seconds=float(health_check.get("timeout_seconds", 2.0)),
failure_threshold=int(health_check.get("failure_threshold", 3)),
recovery_threshold=int(health_check.get("recovery_threshold", 2)),
)
def cmd_init() -> list[CLISubcommand]:
return [RouterServeSubcommand()]
__all__ = ["RouterServeSubcommand", "cmd_init"]
+63 -2
View File
@@ -1,4 +1,65 @@
# SPDX-License-Identifier: Apache-2.0
from fastvideo.entrypoints.streaming.server import run_server
from fastvideo.entrypoints.streaming.server import build_app, run_server
from fastvideo.entrypoints.streaming.session import (
Session,
SessionManager,
SessionState,
)
from fastvideo.entrypoints.streaming.session_store import (
BlobStore,
InMemoryBlobStore,
InMemorySessionStore,
SessionStore,
)
from fastvideo.entrypoints.streaming.gpu_pool import (
GpuPool,
InProcessGpuPool,
PoolAcquireTimeout,
SubprocessGpuPool,
)
from fastvideo.entrypoints.streaming.mock_server import (
MockGenerator,
build_mock_app,
)
from fastvideo.entrypoints.streaming.prompt import (
LLMProvider,
PromptEnhancer,
)
from fastvideo.entrypoints.streaming.prompt.safety import (
PromptSafetyFilter,
SafetyDecision,
)
from fastvideo.entrypoints.streaming.session_logger import (
SessionLogEvent,
SessionLogger,
)
from fastvideo.entrypoints.streaming.stream import (
FragmentedMP4Chunk,
FragmentedMP4Encoder,
)
__all__ = ["run_server"]
__all__ = [
"BlobStore",
"FragmentedMP4Chunk",
"FragmentedMP4Encoder",
"GpuPool",
"InMemoryBlobStore",
"InMemorySessionStore",
"InProcessGpuPool",
"LLMProvider",
"MockGenerator",
"PoolAcquireTimeout",
"PromptEnhancer",
"PromptSafetyFilter",
"SafetyDecision",
"SessionLogEvent",
"SessionLogger",
"build_mock_app",
"Session",
"SessionManager",
"SessionState",
"SessionStore",
"SubprocessGpuPool",
"build_app",
"run_server",
]
+542
View File
@@ -0,0 +1,542 @@
# SPDX-License-Identifier: Apache-2.0
"""GPU pool manager for the streaming server.
Replaces the single-generator path in PR 7.5 with a typed pool
abstraction. Three implementations ship here:
* :class:`InProcessGpuPool` — one in-process ``VideoGenerator``; used
by tests and single-GPU dev deployments.
* :class:`SubprocessGpuPool` — one ``multiprocessing.Process`` per
GPU, each running :func:`worker_main` against a ``GeneratorConfig``.
Jobs are dispatched via ``multiprocessing.Queue``.
* :class:`GpuPool` (abstract) — the interface both use.
Session-to-GPU binding lives in the pool so continuation state stays
on the GPU that generated the previous segment (matching the internal
``gpu_pool.py``'s per-GPU cache behavior). Cross-GPU handoff is
supported via :class:`SessionStore` snapshot + hydrate, which
serializes the state before the migration and rehydrates it on the
new worker.
Typed config: workers start from a :class:`GeneratorConfig` (no flat
LTX-2 kwargs), satisfying the PR 6 + PR 7 contracts that the public
surface doesn't reintroduce the legacy kwarg bag.
"""
from __future__ import annotations
import asyncio
import multiprocessing as mp
import queue
import threading
import time
import uuid
from abc import ABC, abstractmethod
from concurrent.futures import Future
from dataclasses import dataclass, field
from typing import Any, Protocol
from fastvideo.api.schema import (
GeneratorConfig,
GenerationRequest,
GpuPoolConfig,
WarmupConfig,
)
from fastvideo.entrypoints.streaming.session_store import (
InMemorySessionStore,
SessionStore,
)
from fastvideo.entrypoints.streaming.worker import worker_main
from fastvideo.logger import init_logger
logger = init_logger(__name__)
# ---------------------------------------------------------------------------
# Public interface
# ---------------------------------------------------------------------------
class _GeneratorLike(Protocol):
"""Subset the pool calls on a worker-side generator."""
def generate(self, request: GenerationRequest) -> Any:
...
@dataclass
class PoolAssignment:
"""The worker a session is currently bound to."""
gpu_id: int
worker_id: str
pinned_at: float = field(default_factory=time.monotonic)
class GpuPool(ABC):
"""Abstract GPU pool.
``acquire`` binds a session to a worker and holds that binding
across segments so continuation state can stay hot. ``run`` submits
a single ``GenerationRequest`` for a bound session.
Acquire / release are independent of run — a session can run many
segments on one acquired worker, and must release on disconnect.
"""
@abstractmethod
async def acquire(
self,
session_id: str,
*,
timeout: float | None = None,
) -> PoolAssignment:
...
@abstractmethod
async def run(
self,
session_id: str,
request: GenerationRequest,
) -> Any:
...
@abstractmethod
async def release(self, session_id: str) -> None:
...
@abstractmethod
async def shutdown(self) -> None:
...
@abstractmethod
def health(self) -> PoolHealth:
...
@dataclass
class PoolHealth:
total_workers: int
available_workers: int
active_sessions: int
queued_sessions: int = 0
class PoolAcquireTimeout(RuntimeError):
"""Raised when ``acquire`` times out waiting for a free worker."""
# ---------------------------------------------------------------------------
# In-process implementation (single-worker, test / dev)
# ---------------------------------------------------------------------------
class InProcessGpuPool(GpuPool):
"""Single-process pool backed by one :class:`_GeneratorLike`.
This is what PR 7.5's server uses by default; PR 7.6 adds the real
``SubprocessGpuPool`` alternative but keeps this one for tests and
small deployments.
"""
def __init__(
self,
generator: _GeneratorLike,
*,
gpu_id: int = 0,
session_store: SessionStore | None = None,
) -> None:
self._generator = generator
self._gpu_id = gpu_id
self._worker_id = f"inproc-{uuid.uuid4().hex[:6]}"
self._session_store = session_store or InMemorySessionStore()
self._active: dict[str, PoolAssignment] = {}
self._lock = asyncio.Lock()
self._gen_lock = asyncio.Lock()
async def acquire(
self,
session_id: str,
*,
timeout: float | None = None,
) -> PoolAssignment:
async with self._lock:
existing = self._active.get(session_id)
if existing is not None:
return existing
assignment = PoolAssignment(gpu_id=self._gpu_id, worker_id=self._worker_id)
self._active[session_id] = assignment
return assignment
async def run(
self,
session_id: str,
request: GenerationRequest,
) -> Any:
if session_id not in self._active:
raise RuntimeError(f"session {session_id!r} is not acquired on this pool")
# Serialize generator access so one GPU runs one request at a
# time, matching the internal gpu_pool's per-GPU lock.
async with self._gen_lock:
loop = asyncio.get_running_loop()
return await loop.run_in_executor(None, self._generator.generate, request)
async def release(self, session_id: str) -> None:
async with self._lock:
self._active.pop(session_id, None)
async def shutdown(self) -> None:
self._active.clear()
def health(self) -> PoolHealth:
return PoolHealth(
total_workers=1,
available_workers=1 if not self._active else 0,
active_sessions=len(self._active),
)
# ---------------------------------------------------------------------------
# Subprocess implementation (multi-worker, real deployment)
# ---------------------------------------------------------------------------
@dataclass
class _WorkerHandle:
process: Any # mp.Process or compatible handle with is_alive / join / kill
job_queue: mp.Queue
result_queue: mp.Queue
gpu_id: int
worker_id: str
ready: threading.Event
# ``ready`` flips on either successful boot or boot failure so the
# parent stops waiting; ``boot_ok`` is set only on a real ready
# acknowledgement and is what gates pool admission.
boot_ok: threading.Event
shutdown_event: Any # mp.Event is a factory, not a type — Any keeps mypy sane
@dataclass
class _PendingJob:
job_id: str
future: Future
session_id: str
worker_id: str
class SubprocessGpuPool(GpuPool):
"""One ``multiprocessing.Process`` per GPU.
Each worker boots :class:`fastvideo.VideoGenerator` from a typed
:class:`GeneratorConfig` inside the child process (post-
``CUDA_VISIBLE_DEVICES`` setup) and consumes jobs from an mp Queue.
This is the production shape: the parent process stays CPU-only, and
GPU state never crosses process boundaries. Continuation state is
serialized through :class:`SessionStore` for cross-GPU handoff.
PR 7.6 ships this as an opt-in; PR 7.5's in-process pool remains the
default until nightly runs validate the subprocess path.
"""
def __init__(
self,
generator_config: GeneratorConfig,
*,
pool_config: GpuPoolConfig,
warmup_config: WarmupConfig | None = None,
session_store: SessionStore | None = None,
worker_factory: WorkerFactory | None = None,
) -> None:
self._generator_config = generator_config
self._pool_config = pool_config
self._warmup_config = warmup_config or WarmupConfig()
self._session_store = session_store or InMemorySessionStore()
self._worker_factory = worker_factory or _default_worker_factory
self._workers: list[_WorkerHandle] = []
self._available: asyncio.Queue[int] = asyncio.Queue()
self._assignments: dict[str, PoolAssignment] = {}
self._worker_by_id: dict[str, _WorkerHandle] = {}
self._pending: dict[str, _PendingJob] = {}
self._lock = asyncio.Lock()
self._result_reader_tasks: list[asyncio.Task] = []
async def start(self) -> None:
"""Spawn worker processes and wait for each to report ready."""
num_workers = self._pool_config.num_workers or 1
for gpu_id in range(num_workers):
handle = self._worker_factory(
gpu_id=gpu_id,
generator_config=self._generator_config,
warmup_config=self._warmup_config,
)
self._workers.append(handle)
self._worker_by_id[handle.worker_id] = handle
# Wait for each worker's ready event in a thread to avoid
# blocking the event loop.
loop = asyncio.get_running_loop()
await asyncio.gather(*[
loop.run_in_executor(None, handle.ready.wait, self._warmup_config.timeout_seconds)
for handle in self._workers
])
# Start background result readers — one task per worker
# drains its result queue and resolves futures in _pending.
for handle in self._workers:
task = asyncio.create_task(self._drain_results(handle))
self._result_reader_tasks.append(task)
# Only admit workers that successfully booted. Anything that
# failed boot (timeout, crash, error sentinel) stays out of the
# available queue so we never assign a session to it.
for idx, handle in enumerate(self._workers):
if handle.boot_ok.is_set():
await self._available.put(idx)
else:
logger.error(
"pool: worker %s failed to boot; skipping",
handle.worker_id,
)
async def acquire(
self,
session_id: str,
*,
timeout: float | None = None,
) -> PoolAssignment:
async with self._lock:
existing = self._assignments.get(session_id)
if existing is not None:
return existing
try:
idx = await asyncio.wait_for(self._available.get(), timeout=timeout)
except asyncio.TimeoutError as exc:
raise PoolAcquireTimeout(f"no worker available after {timeout}s") from exc
handle = self._workers[idx]
assignment = PoolAssignment(gpu_id=handle.gpu_id, worker_id=handle.worker_id)
async with self._lock:
self._assignments[session_id] = assignment
return assignment
async def run(
self,
session_id: str,
request: GenerationRequest,
) -> Any:
assignment = self._assignments.get(session_id)
if assignment is None:
raise RuntimeError(f"session {session_id!r} not acquired on this pool")
handle = self._worker_by_id[assignment.worker_id]
job_id = uuid.uuid4().hex
future: Future = Future()
self._pending[job_id] = _PendingJob(
job_id=job_id,
future=future,
session_id=session_id,
worker_id=handle.worker_id,
)
# mp.Queue.put can block if the underlying pipe buffer is full;
# offload to a thread so the event loop keeps serving other
# sessions. If the put itself fails, drop the pending entry so
# _drain_results doesn't dangle a future forever.
loop = asyncio.get_running_loop()
try:
await loop.run_in_executor(
None,
handle.job_queue.put,
{
"job_id": job_id,
"request": request
},
)
except Exception:
self._pending.pop(job_id, None)
raise
return await asyncio.wrap_future(future)
async def release(self, session_id: str) -> None:
async with self._lock:
assignment = self._assignments.pop(session_id, None)
if assignment is None:
return
idx = next((i for i, h in enumerate(self._workers) if h.worker_id == assignment.worker_id), None)
if idx is None:
return
# Don't return a dead worker to the pool; otherwise the next
# acquire will hand a session to a process that can't run jobs.
if not self._workers[idx].process.is_alive():
logger.warning(
"pool: worker %s died; not returning to available queue",
self._workers[idx].worker_id,
)
return
await self._available.put(idx)
async def shutdown(self) -> None:
loop = asyncio.get_running_loop()
# Signal all workers in parallel; .put may block on a full pipe,
# so off-load it the same way run() does.
async def _signal(handle: _WorkerHandle) -> None:
try:
handle.shutdown_event.set()
await loop.run_in_executor(None, handle.job_queue.put, None)
except Exception: # pragma: no cover - best-effort cleanup
pass
await asyncio.gather(*(_signal(h) for h in self._workers))
# Join in parallel so total shutdown is bounded by the slowest
# worker, not the sum of all timeouts.
await asyncio.gather(*(loop.run_in_executor(None, handle.process.join, 5.0) for handle in self._workers))
for handle in self._workers:
if handle.process.is_alive():
handle.process.kill()
for task in self._result_reader_tasks:
task.cancel()
self._result_reader_tasks.clear()
self._workers.clear()
self._worker_by_id.clear()
def health(self) -> PoolHealth:
return PoolHealth(
total_workers=len(self._workers),
available_workers=self._available.qsize(),
active_sessions=len(self._assignments),
)
async def _drain_results(self, handle: _WorkerHandle) -> None:
loop = asyncio.get_running_loop()
try:
while not handle.shutdown_event.is_set():
try:
msg = await loop.run_in_executor(None, _safe_queue_get, handle.result_queue, 0.5)
except Exception:
logger.exception("pool: worker %s result reader failed", handle.worker_id)
return
if msg is None:
continue
job_id = msg.get("job_id")
if job_id is None:
continue
pending = self._pending.pop(job_id, None)
if pending is None:
continue
if msg.get("kind") == "error":
pending.future.set_exception(RuntimeError(msg["error"]))
else:
pending.future.set_result(msg.get("result"))
finally:
# If we exit for any reason — shutdown, exception, cancel —
# surface that to any in-flight jobs on this worker so their
# await never hangs on a future no one will resolve.
for jid in [jid for jid, job in self._pending.items() if job.worker_id == handle.worker_id]:
pending = self._pending.pop(jid, None)
if pending is not None and not pending.future.done():
pending.future.set_exception(
RuntimeError(f"worker {handle.worker_id} result reader exited "
"with pending jobs"))
def _safe_queue_get(q: mp.Queue, timeout: float) -> Any | None:
try:
return q.get(timeout=timeout)
except queue.Empty:
return None
# ---------------------------------------------------------------------------
# Worker process
# ---------------------------------------------------------------------------
class WorkerFactory(Protocol):
def __call__(
self,
*,
gpu_id: int,
generator_config: GeneratorConfig,
warmup_config: WarmupConfig,
) -> _WorkerHandle:
...
def _default_worker_factory(
*,
gpu_id: int,
generator_config: GeneratorConfig,
warmup_config: WarmupConfig,
) -> _WorkerHandle:
"""Spawn a real multiprocessing worker.
The child process calls :func:`worker_main` which constructs a
:class:`VideoGenerator` from ``generator_config`` and runs a
blocking job loop. The ``ready`` event flips after the warmup
request completes.
"""
ctx = mp.get_context("spawn")
job_queue: mp.Queue = ctx.Queue()
result_queue: mp.Queue = ctx.Queue()
ready = threading.Event()
boot_ok = threading.Event()
shutdown_event = ctx.Event()
worker_id = f"gpu{gpu_id}-{uuid.uuid4().hex[:6]}"
process = ctx.Process(
target=worker_main,
kwargs={
"gpu_id": gpu_id,
"worker_id": worker_id,
"generator_config": generator_config,
"warmup_config": warmup_config,
"job_queue": job_queue,
"result_queue": result_queue,
"shutdown_event": shutdown_event,
},
daemon=False,
)
process.start()
# Block the parent-side ``ready`` flag until the worker posts a
# ready acknowledgement on the result queue. We drain that single
# sentinel here; subsequent results belong to jobs. ``boot_ok``
# only flips on a real ready; on error we set ``ready`` to unblock
# the parent's wait but leave ``boot_ok`` clear so the pool keeps
# the worker out of the available queue.
def _await_ready() -> None:
while not shutdown_event.is_set():
try:
msg = result_queue.get(timeout=1.0)
except queue.Empty:
continue
if isinstance(msg, dict) and msg.get("kind") == "ready":
boot_ok.set()
ready.set()
return
if isinstance(msg, dict) and msg.get("kind") == "error":
logger.error("pool: worker %s failed to boot: %s", worker_id, msg.get("error"))
ready.set()
return
threading.Thread(target=_await_ready, daemon=True).start()
return _WorkerHandle(
process=process,
job_queue=job_queue,
result_queue=result_queue,
gpu_id=gpu_id,
worker_id=worker_id,
ready=ready,
boot_ok=boot_ok,
shutdown_event=shutdown_event,
)
__all__ = [
"GpuPool",
"InProcessGpuPool",
"PoolAcquireTimeout",
"PoolAssignment",
"PoolHealth",
"SubprocessGpuPool",
"WorkerFactory",
"worker_main",
]
@@ -0,0 +1,122 @@
# SPDX-License-Identifier: Apache-2.0
"""Mock streaming server — a frontend dev aid.
Boots the same FastAPI app the real streaming server uses, but backs
it with :class:`InProcessGpuPool` wrapping a synthetic generator that
emits pre-baked RGB frames. No GPU or model weights required.
Use cases:
* Frontend development without a real model loaded.
* Integration tests that exercise the WS protocol end-to-end.
* Reproducing protocol bugs locally.
Launch: ``python -m fastvideo.entrypoints.streaming.mock_server``.
"""
from __future__ import annotations
import argparse
import time
from dataclasses import dataclass
from typing import Any
import numpy as np
from fastvideo.api.schema import (
ContinuationState,
GenerationRequest,
GeneratorConfig,
SamplingConfig,
ServeConfig,
StreamingConfig,
)
from fastvideo.entrypoints.streaming.server import build_app
@dataclass
class MockGenerator:
"""Generator stand-in that returns synthetic gradient frames.
Each call produces one segment worth of frames whose pixels vary by
a constant derived from the request seed and segment index. Latency
is configurable via ``sleep_ms`` so the caller can exercise slow-
generate scenarios without spinning a GPU.
"""
sleep_ms: float = 0.0
def generate(self, request: GenerationRequest) -> dict[str, Any]:
if self.sleep_ms:
time.sleep(self.sleep_ms / 1000.0)
width = max(16, request.sampling.width)
height = max(16, request.sampling.height)
num_frames = max(1, request.sampling.num_frames)
frames = [_gradient_frame(height, width, idx, seed=request.sampling.seed) for idx in range(num_frames)]
state = ContinuationState(
kind="ltx2.v1",
payload={
"schema_version": 1,
"segment_index": 0,
"source_prompt": request.prompt,
},
)
return {
"frames": frames,
"audio_sample_rate": 24000,
"state": state,
}
def _gradient_frame(height: int, width: int, idx: int, *, seed: int) -> np.ndarray:
base = (idx * 17 + seed * 3) % 256
row = np.linspace(base, (base + 64) % 256, width, dtype=np.uint8)
frame = np.tile(row, (height, 1))
stacked = np.stack([frame, np.roll(frame, 8, axis=1), np.roll(frame, 16, axis=1)], axis=-1)
return stacked.astype(np.uint8)
def build_mock_app(*, sleep_ms: float = 0.0):
"""Build a FastAPI app backed by :class:`MockGenerator`."""
serve_config = ServeConfig(
generator=GeneratorConfig(model_path="/models/mock"),
streaming=StreamingConfig(
session_timeout_seconds=120,
generation_segment_cap=6,
),
)
serve_config.default_request.sampling = SamplingConfig(
num_frames=24,
height=256,
width=256,
fps=24,
num_inference_steps=1,
)
return build_app(serve_config, MockGenerator(sleep_ms=sleep_ms))
def main() -> None: # pragma: no cover - CLI entry
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--host", default="127.0.0.1")
parser.add_argument("--port", type=int, default=8000)
parser.add_argument(
"--sleep-ms",
type=float,
default=0.0,
help="Per-segment artificial latency for testing slow paths",
)
args = parser.parse_args()
import uvicorn
app = build_mock_app(sleep_ms=args.sleep_ms)
uvicorn.run(app, host=args.host, port=args.port)
__all__ = [
"MockGenerator",
"build_mock_app",
"main",
]
if __name__ == "__main__": # pragma: no cover - CLI entry
main()
@@ -0,0 +1,36 @@
# SPDX-License-Identifier: Apache-2.0
"""Prompt pipeline for the streaming server.
* :mod:`providers` — LLM backend abstraction + built-in adapters
* :mod:`enhancer` — provider-agnostic enhance / auto-extend / rewrite
operations on top of the provider layer
All of this is optional; the streaming server runs fine without it
(PR 7.5's skeleton never invokes the enhancer). When the operator
enables ``ServeConfig.streaming.prompt.enabled``, the server routes
each ``session_init_v2`` curated prompt through ``enhance`` before the
first segment.
"""
from fastvideo.entrypoints.streaming.prompt.enhancer import (
PromptEnhancer,
PromptOperation,
)
from fastvideo.entrypoints.streaming.prompt.providers.base import (
LLMMessage,
LLMProvider,
LLMProviderError,
LLMRequest,
LLMResponse,
LLMTimeoutError,
)
__all__ = [
"LLMMessage",
"LLMProvider",
"LLMProviderError",
"LLMRequest",
"LLMResponse",
"LLMTimeoutError",
"PromptEnhancer",
"PromptOperation",
]
@@ -0,0 +1,197 @@
# SPDX-License-Identifier: Apache-2.0
"""Provider-agnostic prompt orchestration for the streaming server.
Three operations the streaming server needs:
* ``enhance`` — polish a user prompt (add cinematic detail, fix syntax)
* ``auto_extend`` — generate a follow-on prompt for loop generation
* ``rewrite`` — rewrite a seed prompt for a user-directed rewrite flow
All three share the same orchestration: pick a provider in priority
order, submit an ``LLMRequest``, fall back to the next provider on
retryable errors, and surface a structured :class:`LLMResponse` back
to the caller.
System prompts are loaded from ``system_prompt_dir`` on construction
and can be hot-reloaded via :meth:`PromptEnhancer.reload_system_prompts`.
The streaming server's management endpoint calls that method in
response to a ``rewrite_seed_prompts_started`` frame.
"""
from __future__ import annotations
import enum
import os
from collections.abc import Sequence
from dataclasses import dataclass, replace
from fastvideo.entrypoints.streaming.prompt.providers.base import (
LLMMessage,
LLMProvider,
LLMProviderError,
LLMRequest,
LLMResponse,
)
from fastvideo.logger import init_logger
logger = init_logger(__name__)
class PromptOperation(enum.Enum):
ENHANCE = "enhance"
AUTO_EXTEND = "auto_extend"
REWRITE = "rewrite"
@dataclass
class _SystemPrompts:
enhance: str
auto_extend: str
rewrite: str
_DEFAULT_SYSTEM_PROMPTS = _SystemPrompts(
enhance=("You are a prompt enhancer for cinematic video generation. Given "
"a user prompt, produce an enhanced prompt that is more vivid, "
"specific, and concrete. Keep the subject intact; add lighting, "
"camera, and motion detail. Reply with just the enhanced prompt."),
auto_extend=("You are a video continuation assistant. Given the current "
"sequence of prompts, produce one new prompt that naturally "
"continues the sequence. Reply with just the next prompt."),
rewrite=("You are a creative prompt rewriter. Given a seed prompt, produce "
"a set of alternative prompts that explore different angles, "
"styles, and moods. Reply with one prompt per line."),
)
class PromptEnhancer:
"""Orchestrates prompt operations across a priority-ordered provider
list with structured fallback + hot-reloadable system prompts.
Usage::
enhancer = PromptEnhancer(
providers=[CerebrasProvider(), GroqProvider()],
model="gpt-oss-120b",
system_prompt_dir="/etc/fastvideo/prompts",
)
response = await enhancer.enhance("a fox running through snow")
"""
def __init__(
self,
*,
providers: Sequence[LLMProvider],
model: str,
timeout_ms: int = 20000,
temperature: float = 0.7,
max_tokens: int | None = 256,
system_prompt_dir: str | None = None,
) -> None:
if not providers:
raise ValueError("PromptEnhancer requires at least one LLMProvider")
self._providers = list(providers)
self._model = model
self._timeout_ms = timeout_ms
self._temperature = temperature
self._max_tokens = max_tokens
self._system_prompt_dir = system_prompt_dir
self._system_prompts = self._load_system_prompts()
@property
def providers(self) -> list[LLMProvider]:
return list(self._providers)
def register_provider(self, provider: LLMProvider, *, priority: int = -1) -> None:
"""Insert an additional provider. ``priority=0`` makes it primary;
``priority=-1`` (default) appends as a fallback."""
if priority < 0:
self._providers.append(provider)
else:
self._providers.insert(priority, provider)
def reload_system_prompts(self) -> None:
"""Re-read the system prompt files from ``system_prompt_dir``.
The streaming server exposes this via a management endpoint so
operators can iterate on prompt templates without restarting
workers.
"""
self._system_prompts = self._load_system_prompts()
logger.info("prompt enhancer: reloaded system prompts from %s", self._system_prompt_dir or "defaults")
async def enhance(self, prompt: str) -> LLMResponse:
return await self._run(
PromptOperation.ENHANCE,
system=self._system_prompts.enhance,
user=prompt,
)
async def auto_extend(self, prior_prompts: Sequence[str]) -> LLMResponse:
user = "\n".join(prior_prompts)
return await self._run(
PromptOperation.AUTO_EXTEND,
system=self._system_prompts.auto_extend,
user=user,
)
async def rewrite(self, seed_prompt: str) -> LLMResponse:
return await self._run(
PromptOperation.REWRITE,
system=self._system_prompts.rewrite,
user=seed_prompt,
)
async def _run(
self,
operation: PromptOperation,
*,
system: str,
user: str,
) -> LLMResponse:
request = LLMRequest(
messages=[
LLMMessage(role="system", content=system),
LLMMessage(role="user", content=user),
],
model=self._model,
max_tokens=self._max_tokens,
temperature=self._temperature,
timeout_ms=self._timeout_ms,
)
last_error: LLMProviderError | None = None
for idx, provider in enumerate(self._providers):
try:
response = await provider.complete(request)
if idx > 0:
# Mark the fallback flag without losing any other
# response fields the provider populated.
response = replace(response, fallback_used=True)
return response
except LLMProviderError as exc:
logger.warning("prompt %s: provider %s failed: %s; trying next", operation.value, provider.name, exc)
last_error = exc
if not exc.retryable:
break
assert last_error is not None
raise last_error
def _load_system_prompts(self) -> _SystemPrompts:
if not self._system_prompt_dir:
return _DEFAULT_SYSTEM_PROMPTS
return _SystemPrompts(
enhance=_read_prompt(self._system_prompt_dir, "enhance.txt", _DEFAULT_SYSTEM_PROMPTS.enhance),
auto_extend=_read_prompt(self._system_prompt_dir, "auto_extend.txt", _DEFAULT_SYSTEM_PROMPTS.auto_extend),
rewrite=_read_prompt(self._system_prompt_dir, "rewrite.txt", _DEFAULT_SYSTEM_PROMPTS.rewrite),
)
def _read_prompt(dirname: str, filename: str, default: str) -> str:
path = os.path.join(dirname, filename)
if not os.path.exists(path):
return default
with open(path, encoding="utf-8") as f:
content = f.read().strip()
return content or default
__all__ = ["PromptEnhancer", "PromptOperation"]
@@ -0,0 +1,24 @@
# SPDX-License-Identifier: Apache-2.0
"""LLM provider implementations used by the prompt enhancer."""
from fastvideo.entrypoints.streaming.prompt.providers.base import (
LLMMessage,
LLMProvider,
LLMProviderError,
LLMRequest,
LLMResponse,
LLMTimeoutError,
)
from fastvideo.entrypoints.streaming.prompt.providers.cerebras import (
CerebrasProvider, )
from fastvideo.entrypoints.streaming.prompt.providers.groq import GroqProvider
__all__ = [
"CerebrasProvider",
"GroqProvider",
"LLMMessage",
"LLMProvider",
"LLMProviderError",
"LLMRequest",
"LLMResponse",
"LLMTimeoutError",
]
@@ -0,0 +1,101 @@
# SPDX-License-Identifier: Apache-2.0
"""Shared HTTP path for OpenAI-compatible ``/chat/completions`` providers.
Cerebras and Groq both expose the OpenAI chat-completions schema, so
the request shape, error mapping, and response decoding are identical
between them. This module centralizes that logic; the per-provider
modules stay thin (just defaults + env var wiring).
"""
from __future__ import annotations
import time
from fastvideo.entrypoints.streaming.prompt.providers.base import (
LLMProviderError,
LLMRequest,
LLMResponse,
LLMTimeoutError,
)
async def complete_openai_compatible(
*,
api_key: str | None,
api_key_hint: str,
base_url: str,
provider_name: str,
request: LLMRequest,
) -> LLMResponse:
"""Issue a chat-completions call and decode the OpenAI response."""
if not api_key:
raise LLMProviderError(
f"{provider_name} provider requires {api_key_hint} "
"(or explicit api_key=...)",
retryable=False,
)
try:
import httpx
except ImportError as exc: # pragma: no cover - optional dep
raise LLMProviderError(
f"{provider_name} provider requires httpx; install httpx",
retryable=False,
) from exc
timeout_s = (request.timeout_ms or 20000) / 1000.0
t0 = time.perf_counter()
try:
async with httpx.AsyncClient(timeout=timeout_s) as client:
response = await client.post(
f"{base_url}/chat/completions",
headers={
"Authorization": f"Bearer {api_key}",
"Content-Type": "application/json",
},
json={
"model": request.model,
"messages": [{
"role": m.role,
"content": m.content
} for m in request.messages],
"max_tokens": request.max_tokens,
"temperature": request.temperature,
},
)
except httpx.TimeoutException as exc:
raise LLMTimeoutError(f"{provider_name} timed out after {timeout_s}s") from exc
except httpx.HTTPError as exc:
raise LLMProviderError(f"{provider_name} HTTP error: {exc}") from exc
if response.status_code >= 400:
# 5xx and 429 (rate-limit) are retryable: another provider may
# succeed. 4xx (auth, bad-request, etc.) are client errors —
# the enhancer should stop fallback traversal.
retryable = (response.status_code >= 500 or response.status_code == 429)
raise LLMProviderError(
f"{provider_name} returned {response.status_code}: "
f"{response.text[:200]}",
retryable=retryable,
)
try:
data = response.json()
except Exception as exc:
# Non-JSON body usually means a proxy / load-balancer error
# page; leave it retryable so a fallback provider can try.
raise LLMProviderError(f"{provider_name} returned non-JSON body: {exc}") from exc
choices = data.get("choices") or []
if not choices:
raise LLMProviderError(f"{provider_name} returned no choices")
content = choices[0].get("message", {}).get("content") or ""
latency_ms = (time.perf_counter() - t0) * 1000.0
return LLMResponse(
content=content.strip(),
provider=provider_name,
model=request.model,
latency_ms=latency_ms,
)
__all__ = ["complete_openai_compatible"]
@@ -0,0 +1,85 @@
# SPDX-License-Identifier: Apache-2.0
"""LLM provider protocol + DTOs used by the prompt enhancer.
Third-party users add a new provider by implementing
:class:`LLMProvider` and registering it with a prompt enhancer
instance. The shipped providers live in sibling modules
(``cerebras.py``, ``groq.py``) and each is ~100-200 LOC — the
provider layer is intentionally thin so the enhancer stays
provider-agnostic.
"""
from __future__ import annotations
from dataclasses import dataclass
from typing import Literal, Protocol, runtime_checkable
@dataclass
class LLMMessage:
role: Literal["system", "user", "assistant"]
content: str
@dataclass
class LLMRequest:
messages: list[LLMMessage]
model: str
max_tokens: int | None = None
temperature: float | None = None
timeout_ms: int | None = None
@dataclass
class LLMResponse:
content: str
provider: str
model: str
latency_ms: float
fallback_used: bool = False
class LLMProviderError(RuntimeError):
"""Raised when an LLM provider fails a request.
``retryable`` controls whether the enhancer falls back to the next
provider. It is settable per-instance so the same exception type
can describe retryable transport errors (5xx, 429) and
non-retryable client errors (4xx auth/bad-request) without forcing
a separate subclass for every status family.
"""
def __init__(self, message: str, *, retryable: bool = True) -> None:
super().__init__(message)
self.retryable = retryable
class LLMTimeoutError(LLMProviderError):
"""Raised when an LLM provider times out — always retryable."""
def __init__(self, message: str) -> None:
super().__init__(message, retryable=True)
@runtime_checkable
class LLMProvider(Protocol):
"""Provider interface every LLM adapter implements.
Providers are async-first because every built-in implementation
talks to an HTTP API. Synchronous providers can wrap their call in
``asyncio.to_thread`` internally.
"""
name: str
async def complete(self, request: LLMRequest) -> LLMResponse:
...
__all__ = [
"LLMMessage",
"LLMProvider",
"LLMProviderError",
"LLMRequest",
"LLMResponse",
"LLMTimeoutError",
]
@@ -0,0 +1,44 @@
# SPDX-License-Identifier: Apache-2.0
"""Cerebras LLM provider (OpenAI-compatible chat endpoint)."""
from __future__ import annotations
import os
from dataclasses import dataclass
from fastvideo.entrypoints.streaming.prompt.providers._openai_compat import (
complete_openai_compatible, )
from fastvideo.entrypoints.streaming.prompt.providers.base import (
LLMRequest,
LLMResponse,
)
_DEFAULT_BASE_URL = "https://api.cerebras.ai/v1"
_API_KEY_ENV = "CEREBRAS_API_KEY"
@dataclass
class CerebrasProvider:
"""Cerebras inference adapter.
``api_key`` falls back to ``CEREBRAS_API_KEY`` when unset.
"""
api_key: str | None = None
base_url: str = _DEFAULT_BASE_URL
name: str = "cerebras"
def __post_init__(self) -> None:
if self.api_key is None:
self.api_key = os.environ.get(_API_KEY_ENV)
async def complete(self, request: LLMRequest) -> LLMResponse:
return await complete_openai_compatible(
api_key=self.api_key,
api_key_hint=_API_KEY_ENV,
base_url=self.base_url,
provider_name=self.name,
request=request,
)
__all__ = ["CerebrasProvider"]
@@ -0,0 +1,46 @@
# SPDX-License-Identifier: Apache-2.0
"""Groq LLM provider (OpenAI-compatible chat endpoint)."""
from __future__ import annotations
import os
from dataclasses import dataclass
from fastvideo.entrypoints.streaming.prompt.providers._openai_compat import (
complete_openai_compatible, )
from fastvideo.entrypoints.streaming.prompt.providers.base import (
LLMRequest,
LLMResponse,
)
_DEFAULT_BASE_URL = "https://api.groq.com/openai/v1"
_API_KEY_ENV = "GROQ_API_KEY"
@dataclass
class GroqProvider:
"""Groq inference adapter.
Identical wire format to :class:`CerebrasProvider`; both go through
:func:`complete_openai_compatible`. The two providers differ only
in base URL, env var, and model id conventions.
"""
api_key: str | None = None
base_url: str = _DEFAULT_BASE_URL
name: str = "groq"
def __post_init__(self) -> None:
if self.api_key is None:
self.api_key = os.environ.get(_API_KEY_ENV)
async def complete(self, request: LLMRequest) -> LLMResponse:
return await complete_openai_compatible(
api_key=self.api_key,
api_key_hint=_API_KEY_ENV,
base_url=self.base_url,
provider_name=self.name,
request=request,
)
__all__ = ["GroqProvider"]
@@ -0,0 +1,82 @@
# SPDX-License-Identifier: Apache-2.0
"""Rewrite payload builder.
The UI's "rewrite seed prompts" flow asks the enhancer to produce a
batch of alternative prompts given one seed. This module packages the
seed + options into the payload the enhancer expects and unpacks the
response back into a typed :class:`RewriteResult`.
Separating this from :mod:`enhancer` keeps the enhancer provider-
agnostic; anything UI-specific (how many alternatives to request, how
to split the response, temperature) lives here.
"""
from __future__ import annotations
import re
from dataclasses import dataclass
from fastvideo.entrypoints.streaming.prompt.enhancer import PromptEnhancer
_LEADING_MARKER_RE = re.compile(r"^(?:[-*•]\s*|\d+\s*[.)]\s*)+")
@dataclass
class RewriteOptions:
count: int = 3
"""Number of alternative prompts to request."""
temperature: float | None = None
@dataclass
class RewriteResult:
seed_prompt: str
alternatives: list[str]
provider: str
model: str
latency_ms: float
fallback_used: bool = False
async def build_rewrite(
enhancer: PromptEnhancer,
seed_prompt: str,
*,
options: RewriteOptions | None = None,
) -> RewriteResult:
"""Run a rewrite op through the enhancer and return a typed result."""
if not seed_prompt.strip():
raise ValueError("rewrite seed prompt must be non-empty")
options = options or RewriteOptions()
response = await enhancer.rewrite(seed_prompt)
alternatives = _split_response(response.content, limit=options.count)
return RewriteResult(
seed_prompt=seed_prompt,
alternatives=alternatives,
provider=response.provider,
model=response.model,
latency_ms=response.latency_ms,
fallback_used=response.fallback_used,
)
def _split_response(content: str, *, limit: int) -> list[str]:
"""Split the LLM response into discrete prompt candidates.
The shipped system prompt instructs the model to emit one prompt
per line; this function is forgiving about numbered lists or
leading bullets so user-supplied system prompts don't break it.
"""
lines = [line.strip() for line in content.splitlines() if line.strip()]
cleaned: list[str] = []
for line in lines:
stripped = _LEADING_MARKER_RE.sub("", line).strip()
if stripped:
cleaned.append(stripped)
return cleaned[:max(1, limit)]
__all__ = [
"RewriteOptions",
"RewriteResult",
"build_rewrite",
]

Some files were not shown because too many files have changed in this diff Show More