* feat(studio): add Tauri native GGUF intake * feat(studio): polish native GGUF intake * [pre-commit.ci] auto fixes from pre-commit.com hooks for more information, see https://pre-commit.ci * fix(studio): load backend helpers during local setup * fix(studio): acquire native load lease before unload * Studio: harden native path lease verification and Tauri intake - Wrap path.resolve(strict=True) and Path.stat() in NativePathLeaseError so a deleted or unmounted GGUF returns 400 instead of leaking the full filesystem path through the generic load_model/validate_model handler. - Re-apply _reject_network_or_device_path to the resolved canonical path for defense in depth after symlink resolution. - Replace try/except ValueError pattern in the device-path guard with Path.is_relative_to; the previous shape silently swallowed NativePathLeaseError (which subclasses ValueError) so /dev,/proc,/sys were never actually rejected. - Broaden the lease redaction regex and dict-key check (Python and Rust diagnostics) to cover both native_path_lease and nativePathLease so the camelCase form emitted by Tauri/frontend payloads is also redacted. - Hoist the redact_native_paths import to module top in loggers/handlers; the recursive filter no longer pays a per-record import lookup. - Persist activeNativePathToken in the chat runtime store so the rollback branch can mint a fresh lease and reload the previous native GGUF when a new load fails after unload; clear it in clearCheckpoint and overwrite it on each successful load. - use-native-drop: read options through a ref so the Tauri onDragDropEvent listener is registered once and stays attached across option changes; reject ambiguous multi-file drops up front instead of silently registering only the first GGUF. - pick_native_model: use an async pick_file with a tokio oneshot channel instead of blocking_pick_file so the Tokio worker is not held for the duration of the OS dialog. - registerNativeModelPath: drop the duplicate sourceKind argument; the Rust command parameter is source_kind. - install_python_stack: insert the script directory (studio/) on sys.path; the previous insert pointed at studio/backend/ which does not satisfy `from backend.utils.wheel_utils import ...`. * install_python_stack: keep _BACKEND_DIR on sys.path Restore the studio/backend insertion. Although the immediately following `from backend.utils.wheel_utils import (...)` is satisfied by studio/ already being on sys.path[0] when invoked as `python studio/install_python_stack.py`, wheel_utils itself runs `from utils.native_path_leases import ...`, which requires studio/backend/ to be importable. Without the backend insertion, the existing tests/python/test_install_python_stack.py collection fails with ModuleNotFoundError: No module named 'utils'. * Studio: tighten native path lease lifecycle and Tauri intake IPC - register_native_model_path now hardcodes NativePathSourceKind::Drop on the Rust side and the frontend stops sending source_kind. The previous JS payload (source_kind only) never reached the Rust deserializer because Tauri's default ArgumentCase::Camel maps the Rust parameter source_kind to the JS key sourceKind, so drag/drop registration silently failed. Hardcoding the source kind also keeps audit metadata trustworthy on this command. - Add native_path_secret_removed_for_child_start context manager and wrap multiprocessing.Process.start() at the inference, export, training, and data-recipe job spawn sites. The previous wrapper-only scrub left UNSLOTH_STUDIO_NATIVE_PATH_LEASE_SECRET visible to spawn-platform import-time worker code. The wrapper run_without_native_path_secret stays as defense-in-depth inside the child. - Stop passing exc_info=True from the native-grant load/validate error logs in routes/inference.py. The structlog filter_sensitive_data processor runs before the renderer, so ConsoleRenderer formatted tracebacks bypassed redaction; the redacted str(e) preserves the message text. - Replace the os.path.normcase string equality on the resolved canonical path with Path.samefile (with a normcase fallback) so Windows leases that differ only in extended-length \\?\ prefix or short-name spelling are accepted. - Wrap consumeNativePathToken in its own try/catch in the chat runtime rollback. If the previous native-model token has aged out of TOKEN_TTL we now surface a clear modelsError instead of silently swallowing the rollback inside the outer catch. - Reject non-ASCII lease strings in _split_lease and convert UnicodeEncodeError / binascii.Error / ValueError raised by _b64decode into NativePathLeaseError so verify_native_path_lease never escapes raw exceptions to the route handler. - Tighten dropStateForPaths to mark multi-file payloads invalid so the overlay matches the post-fix drop handler that rejects the same payload. - Replace the one-shot fetch in useNativePathLeasesSupported with a delayed-retry loop so the picker/drop becomes available once the backend is up rather than staying disabled for the rest of the session after a transient failure. - Drop the unused setActiveNativePathToken setter; the value is set via setState directly in use-chat-model-runtime. - Add a toast on auto-load failure in use-native-drop so a collapsed model selector does not hide the error. - Burn the lease nonce before _validate_current_stat so a stat-failed lease is single-use even if a later state change happens to match the original size/mtime. * Studio: cache lease secret, harden native path stat checks, polish intake UX - Cache the decoded UNSLOTH_STUDIO_NATIVE_PATH_LEASE_SECRET on first verify and validate that it is base64-decodable and at least 32 bytes. Subsequent _decode_secret calls return from the cache and never touch os.environ, so concurrent /api/inference/load and /api/health requests no longer race with native_path_secret_removed_for_child_start scrubbing the env. native_path_leases_supported now wraps _decode_secret so the health flag matches what verify_native_path_lease actually accepts. - Replace path.is_file()/is_dir() + path.stat() with os.lstat() in _validate_current_stat and explicitly reject S_ISLNK; size and mtime checks now refer to the link itself, closing the same-size+same-mtime symlink-swap window that the prior follow-symlink stat() left open. - Add an issued_at_ms < expires_at_ms sanity check in _validate_payload to reject internally inconsistent (HMAC-protected) lease payloads. - Sort _NATIVE_PATH_REDACTIONS by length (descending) before iterating in redact_native_paths so a longer registered path is replaced before a shorter prefix path; otherwise logs containing /foo/X.gguf.bak after only /foo/X.gguf was registered would leak the .bak suffix. - classify_existing_path now re-checks the canonical path with symlink_metadata after canonicalize, so a regular file that is replaced with a symlink in the small canonicalize window is rejected at registration. - ModelSelector renders the local file picker as its own block (not in the eject ternary), so a user with an active model can still replace it via the picker rather than only via drag/drop. - useNativePathLeasesSupported caps the readiness probe at MAX_READINESS_POLLS (60 = ~5 minutes) and aborts the in-flight fetch on unmount via AbortController, so a permanently-disabled backend stops generating sustained traffic and hot-reload no longer leaks open connections. - useChooseNativeModel returns a stable useCallback closure and guards the OS dialog with a useRef so rapid double-clicks cannot open multiple dialogs and orphan Rust tokens. - Branch the multi-file drop toast: if no GGUF was present we say "Only .gguf model files can be dropped here." and otherwise "Drop a single .gguf model file." so users dropping non-GGUF attachments get an accurate explanation. * native_path_leases: lstat the signed canonical path before resolving The earlier change to lstat inside _validate_current_stat operates on grant.canonical_path, which is the post-resolve target. If the user atomically replaces the originally-signed file with a symlink to a different file of identical size and mtime, path.resolve(strict=True) follows the symlink, samefile returns True (both ends share the new inode), and the lstat in _validate_current_stat sees the regular target file rather than the symlink, so the swap goes undetected. Add an os.lstat on the signed canonical path before path.resolve(strict=True), and reject S_ISLNK there. The lstat in _validate_current_stat stays as defense-in-depth for swaps that occur strictly between resolve and stat. * Studio: scrub native lease secret before mp.Queue spawn and tighten lease lifecycle - Move _CTX.Queue / _CTX.Event / _CTX.Process construction inside native_path_secret_removed_for_child_start at the inference, export, training and data-recipe spawn sites. The first Queue creation lazily spawns Python's multiprocessing.resource_tracker child, so when it ran outside the scrub context the tracker process inherited the lease secret. Reproduced via the proc filesystem environ entry; the wrapped order keeps the tracker clean. - native_path_secret_removed_for_child_start now refcounts entries: the env var is popped on the first entry and restored only when the last context exits. Concurrent training/inference/export starts no longer serialize on the env lock across the entire proc.start yield, while still guaranteeing the env stays empty for the duration of every overlapping spawn. - run_without_native_path_secret now also nulls the module-level cached lease secret. With the existing spawn-only multiprocessing context the cache is irrelevant in practice, but a future fork caller would otherwise inherit the in-memory secret even though the env var was scrubbed. - filter_sensitive_data now applies the native lease key check on the top-level event_dict, not only on nested dicts, so a logger call that includes a lease value as a top-level keyword field actually redacts it (the bare value does not match the prefix-anchored regex). - chat-page loadNativeModelIntent now passes intent.id to clearModelIntent so a second drag-drop during an in-flight first auto-load is not wiped from the chip area when the first resolves. - Bump useNativePathLeasesSupported's MAX_READINESS_POLLS from 60 to 720 so first-run installs that compile llama.cpp from source or download large CUDA wheels (well past 5 minutes) don't permanently disable the native picker. * native_path_leases: serialize first-decode against scrub context _decode_secret used a separate _SECRET_INIT_LOCK from the env scrub's _NATIVE_PATH_ENV_LOCK, so the very first decode (before the cache is populated) could race a concurrent native_path_secret_removed_for_child_start and read os.environ during the env-empty window, raising "Native path grants require the managed desktop backend." Subsequent calls hit the cache and were already safe. Acquire _NATIVE_PATH_ENV_LOCK around the env read inside _SECRET_INIT_LOCK and fall back to _SCRUB_SAVED_SECRET when the scrub has temporarily popped the env var. Lock ordering (init then env) is consistent with no other caller, so no deadlock. * Studio: surface native model load errors and harden native path label cache - Native model load and validate now bubble up the actual exception (with paths redacted) and apply the same friendly-error rewrite the non-native path uses, so users see "CUDA OOM", "trust_remote_code required", etc. instead of a generic "Failed to load native model: <label>". - run_without_native_path_secret now also nulls _SCRUB_SAVED_SECRET so a forked grandchild that imports native_path_leases cannot recover the secret via the scrub-aware fallback in _decode_secret. - _NATIVE_PATH_LABELS now has its own 10000-entry cap independent of the 100-entry redaction list, so display_label_for_native_path no longer falls back to returning the raw canonical path after 101 native paths in one session. Redaction list keeps the 100-entry cap for log-scan performance. - _validate_payload now also rejects null bytes in display_label, which is echoed back in HTTP responses and log lines. * Studio: harden native path lease validation and chained native rollback - child_env_without_native_path_secret now copies os.environ under _NATIVE_PATH_ENV_LOCK so a concurrent scrub-context env pop cannot raise RuntimeError: dictionary changed size during iteration in a background hardware scan or other env reader. - _validate_payload and grant construction route every signed numeric field (version, issued_at_ms, expires_at_ms, size_bytes, modified_ms) through new _required_int / _optional_int helpers that wrap raw int() ValueError into NativePathLeaseError. The single upstream catcher produces 400 instead of 500 for malformed signed payloads. - verify_native_path_lease now runs _validate_current_stat before _consume_nonce, so a transient stat error on the canonical path no longer permanently burns the nonce. Concurrent verifies still serialize through _consume_nonce, so single-use is preserved. - Chained native model rollback now restores activeNativePathToken in the chat runtime store after a successful rollback loadModel. Without this, a second consecutive failed switch could not re-roll-back because the store token had been overwritten by the failed attempt. - validate_model now applies the same not_supported_hints friendly rewrite to native model errors that load_model already does, so a native .gguf that fails validation with an upstream "is not supported" message gets the same actionable wording as the non-native branch. * Studio: harden native path log redaction, status disclosure, and chip lifecycle - structlog processor chain now runs format_exc_info before filter_sensitive_data so traceback strings are produced (and then redacted) rather than passed through as untouched (type, value, tb) tuples that the JSON or console renderer formats after the redaction filter has already finished. - native_path_secret_removed_for_child_start clears _CACHED_LEASE_SECRET in addition to popping the env var, so a fork during the scrub window cannot inherit the cached bytes via the parent's heap. Parent verify calls during the window keep working through the existing scrub-aware fallback in _decode_secret. - load_model's except ValueError handler now redacts native paths and uses the native model log label when native_grant_backed is true. Previously a ValueError raised after lease verification (e.g. from ModelConfig.from_identifier or downstream GGUF parsing) returned the raw exception string in the HTTP response body. - llama_cpp_backend now records the native display label at GGUF load time, and /api/inference/status prefers it over the redaction store. After a Python backend restart the redaction store is empty; the attribute keeps the friendly label, and an absolute model_identifier with no other label source falls back to the basename so the canonical path no longer appears in active_model. - reveal_path_token uses native "reveal and select" commands on macOS (open -R) and Windows (explorer /select,) so the file is highlighted in the file manager. Linux keeps the existing parent-directory open. - Native model rollback that fails because the previous token cannot be consumed now throws a rollback-specific Error, and the outer empty catch was replaced with one that re-throws the rollback error. The rollback-specific message now reaches the user instead of being overwritten by the original load error message. - NativeModelChip tracks the Rust token's expiresAtMs on a single setTimeout, disables the Load button at expiry, and relabels it "Select again" with an explanatory tooltip so users do not click into a guaranteed-failure path after the 15-minute TTL elapses. * Studio: tighten native artifact policy, mmproj sibling check, and intake UX - is_open_safe_artifact no longer grants Open for directories. Reveal already handles directory navigation, so the change closes the attack surface where a macOS .app artifact could be launched via open_path_token + open::that_detached. - Display labels are sanitized in classify_existing_path. Control characters in filenames (newlines, tabs, NUL et al.) are replaced with spaces and the label is trimmed and capped, so a file named with embedded newlines cannot inject forged log lines or scramble the UI status panel. - validate_entry_path skips the size_bytes/modified_ms equality check when the operation is Reveal or Open. Cloud-sync agents (Dropbox, iCloud Drive, OneDrive) routinely rewrite extended-attribute metadata which bumps mtime, and the user expects Reveal/Open to remain available for files in synced folders. - llama_cpp_backend gains a _native_grant_backed flag at GGUF load success. /api/inference/status only applies the absolute-path basename fallback when that flag is true, so a non-native absolute local GGUF still reports its canonical model_identifier and unload by identifier keeps working. - Native vision GGUFs now run through _validate_native_mmproj_companion before llama-server starts: the companion mmproj must be a regular file, not a symlink, and must live in the same resolved directory as the granted GGUF. This stops a hostile sibling or symlinked mmproj from being loaded under a single-file lease. - Chained native rollback restructured: the rollback loadModel + state + refresh runs inside its own try/catch that swallows so the outer throw error surfaces the ORIGINAL load failure. The native-token consume-failure case still throws the rollback-specific message early, before the inner block runs, so its actionable guidance is preserved. - Loading-model state and the duplicate-load guard in the chat runtime hook now compare both the model id and the native path token. Two drops or picks with the same basename in different folders no longer silently dedup; the second token is honored. - chat-page loadNativeModelIntent awaits selectModel before clearing the pending intent. If selectModel returns early via dedup or throws, the chip and its token stay so the user can retry instead of losing the selection. - NativeModelChip's Reveal button is disabled when the lease has expired (Rust would reject it anyway), and the Load button label reads "Expired" instead of "Select again" so the disabled element no longer promises an action it cannot perform. * [pre-commit.ci] auto fixes from pre-commit.com hooks for more information, see https://pre-commit.ci --------- Co-authored-by: pre-commit-ci[bot] <66853113+pre-commit-ci[bot]@users.noreply.github.com> Co-authored-by: Daniel Han <danielhanchen@gmail.com>
938 lines
37 KiB
Python
938 lines
37 KiB
Python
# SPDX-License-Identifier: AGPL-3.0-only
|
|
# Copyright 2026-present the Unsloth AI Inc. team. All rights reserved. See /studio/LICENSE.AGPL-3.0
|
|
|
|
"""
|
|
Training backend — subprocess orchestrator.
|
|
|
|
Each training job runs in a fresh subprocess (mp.get_context("spawn")),
|
|
solving the transformers version-switching problem. The old in-process
|
|
UnslothTrainer singleton is only used inside the subprocess (worker.py).
|
|
|
|
This file orchestrates the subprocess lifecycle, pumps events from the
|
|
worker's mp.Queue, and exposes the same API surface to routes/training.py.
|
|
|
|
Pattern follows core/data_recipe/jobs/manager.py.
|
|
"""
|
|
|
|
import json as _json
|
|
import math
|
|
import multiprocessing as mp
|
|
import queue
|
|
import threading
|
|
import time
|
|
import structlog
|
|
from datetime import datetime, timezone
|
|
from loggers import get_logger
|
|
from dataclasses import dataclass, field
|
|
from pathlib import Path
|
|
from typing import Optional, Tuple, Any
|
|
|
|
import matplotlib.pyplot as plt
|
|
from utils.hardware import prepare_gpu_selection
|
|
from utils.native_path_leases import (
|
|
native_path_secret_removed_for_child_start,
|
|
run_without_native_path_secret,
|
|
)
|
|
|
|
logger = get_logger(__name__)
|
|
|
|
_CTX = mp.get_context("spawn")
|
|
|
|
# Plot styling constants
|
|
PLOT_WIDTH = 8
|
|
PLOT_HEIGHT = 3.5
|
|
|
|
|
|
@dataclass
|
|
class TrainingProgress:
|
|
"""Mirror of trainer.TrainingProgress — kept here so the parent process
|
|
never needs to import the heavy ML modules."""
|
|
|
|
epoch: float = 0
|
|
step: int = 0
|
|
total_steps: int = 0
|
|
loss: Optional[float] = None
|
|
learning_rate: Optional[float] = None
|
|
is_training: bool = False
|
|
is_completed: bool = False
|
|
error: Optional[str] = None
|
|
status_message: str = "Ready to train"
|
|
elapsed_seconds: Optional[float] = None
|
|
eta_seconds: Optional[float] = None
|
|
grad_norm: Optional[float] = None
|
|
num_tokens: Optional[int] = None
|
|
eval_loss: Optional[float] = None
|
|
|
|
|
|
class TrainingBackend:
|
|
"""
|
|
Training orchestration backend — subprocess-based.
|
|
Launches a fresh subprocess per training job, communicates via mp.Queue.
|
|
"""
|
|
|
|
FLUSH_THRESHOLD: int = 10
|
|
|
|
def __init__(self):
|
|
# Subprocess state
|
|
self._proc: Optional[mp.Process] = None
|
|
self._event_queue: Any = None
|
|
self._stop_queue: Any = None
|
|
self._pump_thread: Optional[threading.Thread] = None
|
|
self._lock = threading.Lock()
|
|
|
|
# Progress state (updated by pump thread from subprocess events)
|
|
self._progress = TrainingProgress()
|
|
self._should_stop = False
|
|
self._cancel_requested = False # True only for stop(save=False)
|
|
|
|
# Training Metrics (consumed by routes for SSE and /metrics)
|
|
self.loss_history: list = []
|
|
self.lr_history: list = []
|
|
self.step_history: list = []
|
|
self.grad_norm_history: list = []
|
|
self.grad_norm_step_history: list = []
|
|
self.eval_loss_history: list = []
|
|
self.eval_step_history: list = []
|
|
self.eval_enabled: bool = False
|
|
self.current_theme: str = "light"
|
|
|
|
# Job metadata
|
|
self.current_job_id: Optional[str] = None
|
|
self._output_dir: Optional[str] = None
|
|
|
|
# DB persistence
|
|
self._metric_buffer: list[dict] = []
|
|
self._run_finalized: bool = False
|
|
self._db_run_created: bool = False
|
|
self._db_total_steps_set: bool = False
|
|
self._db_config: Optional[dict] = None
|
|
self._db_started_at: Optional[str] = None
|
|
|
|
logger.info("TrainingBackend initialized (subprocess mode)")
|
|
|
|
# ------------------------------------------------------------------
|
|
# Public API (called by routes/training.py)
|
|
# ------------------------------------------------------------------
|
|
|
|
def start_training(self, job_id: str, **kwargs) -> bool:
|
|
"""Spawn a subprocess to run the full training pipeline.
|
|
|
|
All kwargs are serialized into a config dict and sent to the worker.
|
|
Returns True if the subprocess was started successfully.
|
|
"""
|
|
with self._lock:
|
|
if self._proc is not None and self._proc.is_alive():
|
|
logger.warning("Training subprocess already running")
|
|
return False
|
|
|
|
# Join prior pump thread — refuse to start if it won't die
|
|
if self._pump_thread is not None and self._pump_thread.is_alive():
|
|
self._pump_thread.join(timeout = 5.0)
|
|
if self._pump_thread.is_alive():
|
|
logger.warning(
|
|
"Previous pump thread did not exit within 5s — refusing to start"
|
|
)
|
|
return False
|
|
self._pump_thread = None
|
|
|
|
# Build config dict for the subprocess
|
|
config = {
|
|
"model_name": kwargs["model_name"],
|
|
"training_type": kwargs.get("training_type", "LoRA/QLoRA"),
|
|
"hf_token": kwargs.get("hf_token", ""),
|
|
"load_in_4bit": kwargs.get("load_in_4bit", True),
|
|
"max_seq_length": kwargs.get("max_seq_length", 2048),
|
|
"hf_dataset": kwargs.get("hf_dataset", ""),
|
|
"local_datasets": kwargs.get("local_datasets"),
|
|
"local_eval_datasets": kwargs.get("local_eval_datasets"),
|
|
"format_type": kwargs.get("format_type", ""),
|
|
"subset": kwargs.get("subset"),
|
|
"train_split": kwargs.get("train_split", "train"),
|
|
"eval_split": kwargs.get("eval_split"),
|
|
"eval_steps": kwargs.get("eval_steps", 0.00),
|
|
"dataset_slice_start": kwargs.get("dataset_slice_start"),
|
|
"dataset_slice_end": kwargs.get("dataset_slice_end"),
|
|
"custom_format_mapping": kwargs.get("custom_format_mapping"),
|
|
"is_dataset_image": kwargs.get("is_dataset_image", False),
|
|
"is_dataset_audio": kwargs.get("is_dataset_audio", False),
|
|
"is_embedding": kwargs.get("is_embedding", False),
|
|
"num_epochs": kwargs.get("num_epochs", 3),
|
|
"learning_rate": kwargs.get("learning_rate", "2e-4"),
|
|
"batch_size": kwargs.get("batch_size", 2),
|
|
"gradient_accumulation_steps": kwargs.get("gradient_accumulation_steps", 4),
|
|
"warmup_steps": kwargs.get("warmup_steps"),
|
|
"warmup_ratio": kwargs.get("warmup_ratio"),
|
|
"max_steps": kwargs.get("max_steps", 0),
|
|
"save_steps": kwargs.get("save_steps", 0),
|
|
"weight_decay": kwargs.get("weight_decay", 0.001),
|
|
"random_seed": kwargs.get("random_seed", 3407),
|
|
"packing": kwargs.get("packing", False),
|
|
"optim": kwargs.get("optim", "adamw_8bit"),
|
|
"lr_scheduler_type": kwargs.get("lr_scheduler_type", "linear"),
|
|
"use_lora": kwargs.get("use_lora", True),
|
|
"lora_r": kwargs.get("lora_r", 16),
|
|
"lora_alpha": kwargs.get("lora_alpha", 16),
|
|
"lora_dropout": kwargs.get("lora_dropout", 0.0),
|
|
"target_modules": kwargs.get("target_modules"),
|
|
"gradient_checkpointing": kwargs.get("gradient_checkpointing", "unsloth"),
|
|
"use_rslora": kwargs.get("use_rslora", False),
|
|
"use_loftq": kwargs.get("use_loftq", False),
|
|
"train_on_completions": kwargs.get("train_on_completions", False),
|
|
"finetune_vision_layers": kwargs.get("finetune_vision_layers", True),
|
|
"finetune_language_layers": kwargs.get("finetune_language_layers", True),
|
|
"finetune_attention_modules": kwargs.get(
|
|
"finetune_attention_modules", True
|
|
),
|
|
"finetune_mlp_modules": kwargs.get("finetune_mlp_modules", True),
|
|
"enable_wandb": kwargs.get("enable_wandb", False),
|
|
"wandb_token": kwargs.get("wandb_token"),
|
|
"wandb_project": kwargs.get("wandb_project", "unsloth-training"),
|
|
"enable_tensorboard": kwargs.get("enable_tensorboard", False),
|
|
"tensorboard_dir": kwargs.get("tensorboard_dir", "runs"),
|
|
"resume_from_checkpoint": kwargs.get("resume_from_checkpoint"),
|
|
"trust_remote_code": kwargs.get("trust_remote_code", False),
|
|
"gpu_ids": kwargs.get("gpu_ids"),
|
|
}
|
|
|
|
# Derive load_in_4bit from training_type
|
|
if config["training_type"] != "LoRA/QLoRA":
|
|
config["load_in_4bit"] = False
|
|
|
|
# Spawn subprocess — use locals so state is untouched on failure
|
|
resolved_gpu_ids, gpu_selection = prepare_gpu_selection(
|
|
kwargs.get("gpu_ids"),
|
|
model_name = config["model_name"],
|
|
hf_token = config["hf_token"] or None,
|
|
training_type = config["training_type"],
|
|
load_in_4bit = config["load_in_4bit"],
|
|
batch_size = config.get("batch_size", 4),
|
|
max_seq_length = config.get("max_seq_length", 2048),
|
|
lora_rank = config.get("lora_r", 16),
|
|
target_modules = config.get("target_modules"),
|
|
gradient_checkpointing = config.get("gradient_checkpointing", "unsloth"),
|
|
optimizer = config.get("optim", "adamw_8bit"),
|
|
)
|
|
config["resolved_gpu_ids"] = resolved_gpu_ids
|
|
config["gpu_selection"] = gpu_selection
|
|
|
|
from .worker import run_training_process
|
|
|
|
try:
|
|
with native_path_secret_removed_for_child_start():
|
|
event_queue = _CTX.Queue()
|
|
stop_queue = _CTX.Queue()
|
|
|
|
proc = _CTX.Process(
|
|
target = run_without_native_path_secret,
|
|
args = (run_training_process,),
|
|
kwargs = {
|
|
"event_queue": event_queue,
|
|
"stop_queue": stop_queue,
|
|
"config": config,
|
|
},
|
|
daemon = True,
|
|
)
|
|
proc.start()
|
|
except Exception:
|
|
logger.error("Failed to start training subprocess", exc_info = True)
|
|
return False
|
|
|
|
logger.info("Training subprocess started (pid=%s)", proc.pid)
|
|
|
|
# Reset state — safe because old pump thread is confirmed dead
|
|
# and proc.start() succeeded
|
|
self.current_job_id = job_id
|
|
self._should_stop = False
|
|
self._cancel_requested = False
|
|
self._progress = TrainingProgress(
|
|
is_training = True, status_message = "Initializing training..."
|
|
)
|
|
self.loss_history.clear()
|
|
self.lr_history.clear()
|
|
self.step_history.clear()
|
|
self.grad_norm_history.clear()
|
|
self.grad_norm_step_history.clear()
|
|
self.eval_loss_history.clear()
|
|
self.eval_step_history.clear()
|
|
self.eval_enabled = False
|
|
self._output_dir = None
|
|
self._metric_buffer.clear()
|
|
self._run_finalized = False
|
|
self._db_run_created = False
|
|
self._db_total_steps_set = False
|
|
self._db_config = {
|
|
k: v for k, v in config.items() if k not in {"hf_token", "wandb_token"}
|
|
}
|
|
self._db_started_at = datetime.now(timezone.utc).isoformat()
|
|
|
|
# Assign subprocess handles after state reset
|
|
self._event_queue = event_queue
|
|
self._stop_queue = stop_queue
|
|
self._proc = proc
|
|
|
|
# Eagerly create DB run row so the run appears in history during model loading
|
|
self._ensure_db_run_created()
|
|
|
|
# Start event pump thread
|
|
self._pump_thread = threading.Thread(target = self._pump_loop, daemon = True)
|
|
self._pump_thread.start()
|
|
|
|
return True
|
|
|
|
def stop_training(self, save: bool = True) -> bool:
|
|
"""Send stop signal to the training subprocess."""
|
|
self._should_stop = True
|
|
if not save:
|
|
self._cancel_requested = True
|
|
with self._lock:
|
|
if self._stop_queue is not None:
|
|
try:
|
|
self._stop_queue.put({"type": "stop", "save": save})
|
|
except (OSError, ValueError):
|
|
pass
|
|
# Update progress immediately for responsive UI
|
|
self._progress.status_message = (
|
|
"Stopping training and saving checkpoint..."
|
|
if save
|
|
else "Cancelling training..."
|
|
)
|
|
return True
|
|
|
|
def force_terminate(self) -> None:
|
|
"""Force-kill the training subprocess so state can be reset immediately."""
|
|
with self._lock:
|
|
if self._proc is not None and self._proc.is_alive():
|
|
logger.info(
|
|
"Force-terminating training subprocess (pid=%s)", self._proc.pid
|
|
)
|
|
self._proc.terminate()
|
|
proc = self._proc
|
|
|
|
if proc is not None:
|
|
proc.join(timeout = 5.0)
|
|
if proc.is_alive():
|
|
proc.kill()
|
|
proc.join(timeout = 2.0)
|
|
|
|
# Wait for pump thread to finish DB finalization before returning
|
|
# (8s covers SQLite's default 5s lock timeout plus execution overhead)
|
|
if self._pump_thread is not None and self._pump_thread.is_alive():
|
|
self._pump_thread.join(timeout = 8.0)
|
|
|
|
def is_training_active(self) -> bool:
|
|
"""Check if training is currently active."""
|
|
with self._lock:
|
|
# Subprocess alive = active
|
|
if self._proc is not None and self._proc.is_alive():
|
|
return True
|
|
|
|
# Stop was requested and process exited → inactive
|
|
if self._should_stop:
|
|
return False
|
|
|
|
# Check progress state
|
|
p = self._progress
|
|
if p.is_training:
|
|
return True
|
|
if p.is_completed or p.error:
|
|
return False
|
|
|
|
# Check status message for activity indicators
|
|
status_lower = (p.status_message or "").lower()
|
|
if any(
|
|
k in status_lower
|
|
for k in [
|
|
"cancelled",
|
|
"canceled",
|
|
"stopped",
|
|
"completed",
|
|
"ready to train",
|
|
]
|
|
):
|
|
return False
|
|
if any(
|
|
k in status_lower
|
|
for k in [
|
|
"loading",
|
|
"preparing",
|
|
"training",
|
|
"configuring",
|
|
"tokenizing",
|
|
"starting",
|
|
"importing",
|
|
]
|
|
):
|
|
return True
|
|
|
|
return False
|
|
|
|
def get_training_status(self, theme: str = "light") -> Tuple:
|
|
"""Get current training status and loss plot."""
|
|
with self._lock:
|
|
progress = self._progress
|
|
|
|
if not (progress.is_training or progress.is_completed or progress.error):
|
|
return (None, progress)
|
|
|
|
plot = self._create_loss_plot(progress, theme)
|
|
return (plot, progress)
|
|
|
|
def refresh_plot_for_theme(self, theme: str) -> Optional[plt.Figure]:
|
|
"""Refresh plot with new theme."""
|
|
if theme and isinstance(theme, str) and theme in ["light", "dark"]:
|
|
self.current_theme = theme
|
|
if self.loss_history:
|
|
with self._lock:
|
|
progress = self._progress
|
|
return self._create_loss_plot(progress, self.current_theme)
|
|
return None
|
|
|
|
# ------------------------------------------------------------------
|
|
# Compatibility shims — routes/training.py accesses these
|
|
# ------------------------------------------------------------------
|
|
|
|
class _TrainerShim:
|
|
"""Minimal shim so routes that access backend.trainer.* still work."""
|
|
|
|
def __init__(self, backend: "TrainingBackend"):
|
|
self._backend = backend
|
|
self.should_stop = False
|
|
|
|
@property
|
|
def training_progress(self):
|
|
return self._backend._progress
|
|
|
|
@training_progress.setter
|
|
def training_progress(self, value):
|
|
self._backend._progress = value
|
|
|
|
def get_training_progress(self):
|
|
return self._backend._progress
|
|
|
|
def _update_progress(self, **kwargs):
|
|
with self._backend._lock:
|
|
for key, value in kwargs.items():
|
|
if hasattr(self._backend._progress, key):
|
|
setattr(self._backend._progress, key, value)
|
|
|
|
@property
|
|
def trainer(self):
|
|
"""Compatibility shim for routes that access backend.trainer.*"""
|
|
return self._TrainerShim(self)
|
|
|
|
# ------------------------------------------------------------------
|
|
# Event pump (background thread)
|
|
# ------------------------------------------------------------------
|
|
|
|
def _pump_loop(self) -> None:
|
|
"""Background thread: consume events from subprocess → update state."""
|
|
while True:
|
|
if self._proc is None or self._event_queue is None:
|
|
return
|
|
|
|
# Try to read an event
|
|
event = self._read_queue(self._event_queue, timeout_sec = 0.25)
|
|
if event is not None:
|
|
self._handle_event(event)
|
|
continue
|
|
|
|
# No event — check if process is still alive
|
|
if self._proc.is_alive():
|
|
continue
|
|
|
|
# Process exited — drain remaining events
|
|
for e in self._drain_queue(self._event_queue):
|
|
self._handle_event(e)
|
|
|
|
# Mark as done if no explicit complete/error was received
|
|
with self._lock:
|
|
if self._progress.is_training:
|
|
if self._should_stop:
|
|
self._progress.is_training = False
|
|
self._progress.status_message = "Training stopped."
|
|
else:
|
|
self._progress.is_training = False
|
|
self._progress.error = (
|
|
self._progress.error
|
|
or "Training process exited unexpectedly"
|
|
)
|
|
|
|
self._ensure_db_run_created()
|
|
self._finalize_run_in_db(
|
|
status = "stopped" if self._should_stop else "error",
|
|
error_message = None
|
|
if self._should_stop
|
|
else "Training process terminated unexpectedly",
|
|
)
|
|
return
|
|
|
|
def _handle_event(self, event: dict) -> None:
|
|
"""Apply a subprocess event to local state.
|
|
|
|
State updates happen inside self._lock; DB I/O happens after
|
|
releasing it so status-polling API endpoints are never blocked
|
|
by slow SQLite writes.
|
|
"""
|
|
etype = event.get("type")
|
|
db_action: Optional[str] = None
|
|
db_action_kwargs: dict = {}
|
|
|
|
with self._lock:
|
|
if etype == "progress":
|
|
self._progress.step = event.get("step", self._progress.step)
|
|
self._progress.epoch = event.get("epoch", self._progress.epoch)
|
|
# loss/lr are sanitized below; update progress after coercion
|
|
_raw_loss = event.get("loss")
|
|
_raw_lr = event.get("learning_rate")
|
|
try:
|
|
_safe_loss = float(_raw_loss) if _raw_loss is not None else None
|
|
except (TypeError, ValueError):
|
|
logger.debug("Could not convert loss to float: %s", _raw_loss)
|
|
_safe_loss = None
|
|
if _safe_loss is not None and not math.isfinite(_safe_loss):
|
|
_safe_loss = None
|
|
try:
|
|
_safe_lr = float(_raw_lr) if _raw_lr is not None else None
|
|
except (TypeError, ValueError):
|
|
logger.debug(
|
|
"Could not convert learning_rate to float: %s", _raw_lr
|
|
)
|
|
_safe_lr = None
|
|
if _safe_lr is not None and not math.isfinite(_safe_lr):
|
|
_safe_lr = None
|
|
if _safe_loss is not None:
|
|
self._progress.loss = _safe_loss
|
|
if _safe_lr is not None:
|
|
self._progress.learning_rate = _safe_lr
|
|
self._progress.total_steps = event.get(
|
|
"total_steps", self._progress.total_steps
|
|
)
|
|
self._progress.elapsed_seconds = event.get("elapsed_seconds")
|
|
self._progress.eta_seconds = event.get("eta_seconds")
|
|
self._progress.grad_norm = event.get("grad_norm")
|
|
self._progress.num_tokens = event.get("num_tokens")
|
|
self._progress.eval_loss = event.get("eval_loss")
|
|
self._progress.is_training = True
|
|
status = event.get("status_message", "")
|
|
if status:
|
|
self._progress.status_message = status
|
|
|
|
# Update metric histories — reuse sanitized values from above
|
|
step = event.get("step", 0)
|
|
loss = _safe_loss
|
|
lr = _safe_lr
|
|
if step > 0 and loss is not None:
|
|
self.loss_history.append(loss)
|
|
self.lr_history.append(lr if lr is not None else 0.0)
|
|
self.step_history.append(step)
|
|
|
|
grad_norm = event.get("grad_norm")
|
|
gn = None
|
|
if grad_norm is not None:
|
|
try:
|
|
gn = float(grad_norm)
|
|
except (TypeError, ValueError):
|
|
gn = None
|
|
if step > 0 and gn is not None and math.isfinite(gn):
|
|
self.grad_norm_history.append(gn)
|
|
self.grad_norm_step_history.append(step)
|
|
else:
|
|
gn = None
|
|
|
|
eval_loss = event.get("eval_loss")
|
|
if eval_loss is not None:
|
|
try:
|
|
eval_loss = float(eval_loss)
|
|
except (TypeError, ValueError):
|
|
logger.debug(
|
|
"Could not convert eval_loss to float: %s", eval_loss
|
|
)
|
|
eval_loss = None
|
|
if step > 0 and eval_loss is not None and math.isfinite(eval_loss):
|
|
self.eval_loss_history.append(eval_loss)
|
|
self.eval_step_history.append(step)
|
|
self.eval_enabled = True
|
|
else:
|
|
eval_loss = None
|
|
|
|
# Buffer metric for DB flush (loss/lr already sanitized above)
|
|
self._metric_buffer.append(
|
|
{
|
|
"step": step,
|
|
"loss": loss,
|
|
"learning_rate": lr,
|
|
"grad_norm": gn,
|
|
"eval_loss": eval_loss,
|
|
"epoch": event.get("epoch"),
|
|
"num_tokens": event.get("num_tokens"),
|
|
"elapsed_seconds": event.get("elapsed_seconds"),
|
|
}
|
|
)
|
|
|
|
# Decide which DB action to take after releasing the lock
|
|
if not self._db_run_created and self.current_job_id and self._db_config:
|
|
db_action = "create_run"
|
|
db_action_kwargs = {
|
|
"job_id": self.current_job_id,
|
|
"model_name": self._db_config["model_name"],
|
|
"dataset_name": self._db_config.get("hf_dataset")
|
|
or next(
|
|
iter(self._db_config.get("local_datasets") or []), "unknown"
|
|
),
|
|
"config_json": _json.dumps(self._db_config),
|
|
"started_at": self._db_started_at
|
|
or datetime.now(timezone.utc).isoformat(),
|
|
"total_steps": event.get("total_steps"),
|
|
}
|
|
elif (
|
|
event.get("total_steps")
|
|
and self._db_run_created
|
|
and not self._db_total_steps_set
|
|
):
|
|
db_action = "update_total_steps"
|
|
db_action_kwargs = {
|
|
"job_id": self.current_job_id,
|
|
"total_steps": event["total_steps"],
|
|
}
|
|
elif len(self._metric_buffer) >= self.FLUSH_THRESHOLD:
|
|
db_action = "flush"
|
|
|
|
elif etype == "eval_configured":
|
|
self.eval_enabled = True
|
|
|
|
elif etype == "status":
|
|
self._progress.status_message = event.get("message", "")
|
|
self._progress.is_training = True
|
|
|
|
elif etype == "complete":
|
|
self._progress.is_training = False
|
|
self._progress.is_completed = True
|
|
self._output_dir = event.get("output_dir")
|
|
msg = event.get("status_message", "Training completed")
|
|
self._progress.status_message = msg
|
|
if not self._db_run_created and self.current_job_id and self._db_config:
|
|
db_action = "create_and_finalize"
|
|
else:
|
|
db_action = "finalize"
|
|
db_action_kwargs = {
|
|
"status": "stopped" if self._should_stop else "completed",
|
|
"output_dir": self._output_dir,
|
|
}
|
|
|
|
elif etype == "error":
|
|
self._progress.is_training = False
|
|
self._progress.error = event.get("error", "Unknown error")
|
|
logger.error("Training error: %s", event.get("error"))
|
|
stack = event.get("stack", "")
|
|
if stack:
|
|
logger.error("Stack trace:\n%s", stack)
|
|
if not self._db_run_created and self.current_job_id and self._db_config:
|
|
db_action = "create_and_finalize"
|
|
else:
|
|
db_action = "finalize"
|
|
db_action_kwargs = {
|
|
"status": "stopped" if self._should_stop else "error",
|
|
"error_message": event.get("error", "Unknown error"),
|
|
}
|
|
|
|
# --- DB I/O outside the lock ---
|
|
if db_action == "create_run":
|
|
try:
|
|
from storage.studio_db import create_run
|
|
|
|
create_run(
|
|
id = db_action_kwargs["job_id"],
|
|
model_name = db_action_kwargs["model_name"],
|
|
dataset_name = db_action_kwargs["dataset_name"],
|
|
config_json = db_action_kwargs["config_json"],
|
|
started_at = db_action_kwargs["started_at"],
|
|
total_steps = db_action_kwargs["total_steps"],
|
|
)
|
|
self._db_run_created = True
|
|
if db_action_kwargs["total_steps"]:
|
|
self._db_total_steps_set = True
|
|
except Exception:
|
|
logger.warning("Failed to create DB run record", exc_info = True)
|
|
elif db_action == "create_and_finalize":
|
|
self._ensure_db_run_created()
|
|
self._finalize_run_in_db(**db_action_kwargs)
|
|
elif db_action == "update_total_steps":
|
|
try:
|
|
from storage.studio_db import update_run_total_steps
|
|
|
|
update_run_total_steps(
|
|
db_action_kwargs["job_id"], db_action_kwargs["total_steps"]
|
|
)
|
|
self._db_total_steps_set = True
|
|
except Exception:
|
|
logger.warning("Failed to update total_steps in DB", exc_info = True)
|
|
elif db_action == "flush":
|
|
self._flush_metrics_to_db()
|
|
elif db_action == "finalize":
|
|
self._finalize_run_in_db(**db_action_kwargs)
|
|
|
|
def _ensure_db_run_created(self) -> None:
|
|
"""Create the DB row if it doesn't exist yet. Called outside the lock."""
|
|
if self._db_run_created or not self.current_job_id or not self._db_config:
|
|
return
|
|
try:
|
|
from storage.studio_db import create_run
|
|
|
|
dataset_name = self._db_config.get("hf_dataset") or next(
|
|
iter(self._db_config.get("local_datasets") or []), "unknown"
|
|
)
|
|
create_run(
|
|
id = self.current_job_id,
|
|
model_name = self._db_config["model_name"],
|
|
dataset_name = dataset_name,
|
|
config_json = _json.dumps(self._db_config),
|
|
started_at = self._db_started_at
|
|
or datetime.now(timezone.utc).isoformat(),
|
|
total_steps = self._progress.total_steps or None,
|
|
)
|
|
self._db_run_created = True
|
|
except Exception:
|
|
logger.warning(
|
|
"Failed to create DB run record for early failure", exc_info = True
|
|
)
|
|
|
|
def _finalize_run_in_db(
|
|
self,
|
|
status: str,
|
|
error_message: Optional[str] = None,
|
|
output_dir: Optional[str] = None,
|
|
) -> None:
|
|
"""Flush remaining metrics and mark a run as finished in the DB."""
|
|
if not self.current_job_id or not self._db_run_created or self._run_finalized:
|
|
return
|
|
self._flush_metrics_to_db()
|
|
try:
|
|
from storage.studio_db import finish_run
|
|
from utils.downsample import downsample
|
|
|
|
sparkline = downsample(self.loss_history, 50)
|
|
finish_run(
|
|
id = self.current_job_id,
|
|
status = status,
|
|
ended_at = datetime.now(timezone.utc).isoformat(),
|
|
final_step = self._progress.step,
|
|
final_loss = self._progress.loss
|
|
if (
|
|
self._progress.loss is not None
|
|
and math.isfinite(self._progress.loss)
|
|
)
|
|
else None,
|
|
duration_seconds = self._progress.elapsed_seconds,
|
|
loss_sparkline = _json.dumps(sparkline),
|
|
output_dir = output_dir,
|
|
error_message = error_message,
|
|
)
|
|
self._run_finalized = True
|
|
except Exception:
|
|
logger.warning(
|
|
"Failed to finalize run in DB (status=%s)", status, exc_info = True
|
|
)
|
|
|
|
def _flush_metrics_to_db(self) -> None:
|
|
"""Flush buffered metrics to the database and update live progress."""
|
|
if (
|
|
not self._metric_buffer
|
|
or not self.current_job_id
|
|
or not self._db_run_created
|
|
):
|
|
return
|
|
# Cap buffer to prevent unbounded memory growth
|
|
if len(self._metric_buffer) > 500:
|
|
logger.warning(
|
|
"Metric buffer exceeded 500 entries (%d) — trimming oldest",
|
|
len(self._metric_buffer),
|
|
)
|
|
self._metric_buffer = self._metric_buffer[-500:]
|
|
# Snapshot before insert so metrics arriving during the write are preserved
|
|
batch = list(self._metric_buffer)
|
|
try:
|
|
from storage.studio_db import insert_metrics_batch, update_run_progress
|
|
|
|
insert_metrics_batch(self.current_job_id, batch)
|
|
del self._metric_buffer[: len(batch)]
|
|
update_run_progress(
|
|
id = self.current_job_id,
|
|
step = self._progress.step,
|
|
loss = self._progress.loss
|
|
if (
|
|
self._progress.loss is not None
|
|
and math.isfinite(self._progress.loss)
|
|
)
|
|
else None,
|
|
duration_seconds = self._progress.elapsed_seconds,
|
|
)
|
|
except Exception:
|
|
# Leave buffer intact for retry on next flush
|
|
logger.warning("Failed to flush metrics to DB", exc_info = True)
|
|
|
|
@staticmethod
|
|
def _read_queue(q: Any, timeout_sec: float) -> Optional[dict]:
|
|
try:
|
|
return q.get(timeout = timeout_sec)
|
|
except queue.Empty:
|
|
return None
|
|
except (EOFError, OSError, ValueError):
|
|
return None
|
|
|
|
@staticmethod
|
|
def _drain_queue(q: Any) -> list:
|
|
events = []
|
|
while True:
|
|
try:
|
|
events.append(q.get_nowait())
|
|
except queue.Empty:
|
|
return events
|
|
except (EOFError, OSError, ValueError):
|
|
return events
|
|
|
|
# ------------------------------------------------------------------
|
|
# Plot generation (unchanged from original)
|
|
# ------------------------------------------------------------------
|
|
|
|
def _create_loss_plot(
|
|
self, progress: TrainingProgress, theme: str = "light"
|
|
) -> plt.Figure:
|
|
"""Create training loss plot with theme-aware styling."""
|
|
plt.close("all")
|
|
|
|
LIGHT_STYLE = {
|
|
"facecolor": "#ffffff",
|
|
"grid_color": "#d1d5db",
|
|
"line": "#16b88a",
|
|
"text": "#1f2937",
|
|
"empty_text": "#6b7280",
|
|
}
|
|
DARK_STYLE = {
|
|
"facecolor": "#292929",
|
|
"grid_color": "#404040",
|
|
"line": "#4ade80",
|
|
"text": "#e5e7eb",
|
|
"empty_text": "#9ca3af",
|
|
}
|
|
|
|
style = LIGHT_STYLE if theme == "light" else DARK_STYLE
|
|
|
|
fig, ax = plt.subplots(figsize = (PLOT_WIDTH, PLOT_HEIGHT))
|
|
fig.patch.set_facecolor(style["facecolor"])
|
|
ax.set_facecolor(style["facecolor"])
|
|
|
|
if self.loss_history:
|
|
steps = self.step_history
|
|
losses = self.loss_history
|
|
scatter_color = "#60a5fa"
|
|
ax.scatter(
|
|
steps,
|
|
losses,
|
|
s = 16,
|
|
alpha = 0.6,
|
|
color = scatter_color,
|
|
linewidths = 0,
|
|
label = "Training Loss (raw)",
|
|
)
|
|
|
|
MA_WINDOW = 20
|
|
window = min(MA_WINDOW, len(losses))
|
|
|
|
if window >= 2:
|
|
cumsum = [0.0]
|
|
for v in losses:
|
|
cumsum.append(cumsum[-1] + float(v))
|
|
|
|
ma = []
|
|
for i in range(len(losses)):
|
|
start = max(0, i - window + 1)
|
|
denom = i - start + 1
|
|
ma.append((cumsum[i + 1] - cumsum[start]) / denom)
|
|
|
|
ax.plot(
|
|
steps,
|
|
ma,
|
|
color = style["line"],
|
|
linewidth = 2.5,
|
|
alpha = 0.95,
|
|
label = f"Moving Avg ({ma[-1]:.4f})",
|
|
)
|
|
|
|
leg = ax.legend(frameon = False, fontsize = 9)
|
|
for t in leg.get_texts():
|
|
t.set_color(style["text"])
|
|
|
|
ax.set_xlabel("Steps", fontsize = 10, color = style["text"])
|
|
ax.set_ylabel("Loss", fontsize = 10, color = style["text"])
|
|
|
|
if progress.error:
|
|
title = f"Error: {progress.error}"
|
|
elif progress.is_completed:
|
|
loss_str = f"{progress.loss:.4f}" if progress.loss is not None else "--"
|
|
title = f"Training completed! Final loss: {loss_str}"
|
|
elif progress.status_message:
|
|
title = progress.status_message
|
|
elif progress.step > 0:
|
|
loss_str = f"{progress.loss:.4f}" if progress.loss is not None else "--"
|
|
title = f"Epoch: {progress.epoch} | Step: {progress.step}/{progress.total_steps} | Loss: {loss_str}"
|
|
else:
|
|
title = "Training Loss"
|
|
|
|
ax.set_title(
|
|
title, fontsize = 11, fontweight = "bold", pad = 10, color = style["text"]
|
|
)
|
|
ax.grid(True, alpha = 0.4, linestyle = "--", color = style["grid_color"])
|
|
ax.tick_params(colors = style["text"], which = "both")
|
|
ax.spines["top"].set_visible(False)
|
|
ax.spines["right"].set_visible(False)
|
|
ax.spines["bottom"].set_color(style["text"])
|
|
ax.spines["left"].set_color(style["text"])
|
|
else:
|
|
display_msg = (
|
|
progress.status_message
|
|
if progress.status_message
|
|
else "Waiting for training data..."
|
|
)
|
|
ax.text(
|
|
0.5,
|
|
0.5,
|
|
display_msg,
|
|
ha = "center",
|
|
va = "center",
|
|
fontsize = 16,
|
|
color = style["empty_text"],
|
|
transform = ax.transAxes,
|
|
)
|
|
ax.set_xticks([])
|
|
ax.set_yticks([])
|
|
for spine in ax.spines.values():
|
|
spine.set_visible(False)
|
|
|
|
fig.tight_layout()
|
|
return fig
|
|
|
|
def _transfer_to_inference_backend(self) -> bool:
|
|
"""Transfer model to inference backend.
|
|
|
|
With subprocess-based training, the model lives in the subprocess
|
|
and is freed when it exits. Inference must load from the saved
|
|
checkpoint on disk. This is a no-op placeholder.
|
|
"""
|
|
logger.info(
|
|
"_transfer_to_inference_backend: subprocess training — "
|
|
"model must be loaded from disk (output_dir=%s)",
|
|
self._output_dir,
|
|
)
|
|
return False
|
|
|
|
|
|
# ========== GLOBAL INSTANCE ==========
|
|
_training_backend = None
|
|
|
|
|
|
def get_training_backend() -> TrainingBackend:
|
|
"""Get global training backend instance"""
|
|
global _training_backend
|
|
if _training_backend is None:
|
|
_training_backend = TrainingBackend()
|
|
return _training_backend
|