From 0a66dbbdb57a0e2590ab1f7b39b41eed7baf9a51 Mon Sep 17 00:00:00 2001 From: simba Date: Sat, 26 Sep 2026 14:27:42 -0700 Subject: [PATCH 1/7] fix(ssd): reuse only finished prefetches; directory load after the MoE check Follow-ups from the #196 review: - The prefetch reuse check requires a marker that is written only after snapshot returns and the copy has every shard plus the tokenizer files. A download interrupted after the last shard (tokenizer or template missing) is fetched again instead of failing on every later start. - After the snapshot, re-validate and throw ModelDownloadIncomplete (model_load_failed) when the copy is partial; offline, snapshot returns the directory even then. - Switch the main load to the streaming directory only after the non-MoE check, and only when tokenizer files are present, so dense models started with --stream-experts and incomplete copies keep the id-based load that can fill in missing files. - ProgressTracker.finish() draws a last frame at the real fraction before ending the line, so the bar no longer stops at ~90%. - README: give the broadcast_shapes crash as an example, not a fixed top-k. Co-Authored-By: Claude Opus 5.5 --- README.md | 2 +- Sources/SwiftLM/Server.swift | 110 ++++++++++++++++++++++++----------- 2 files changed, 77 insertions(+), 35 deletions(-) diff --git a/README.md b/README.md index bf31a2c..fecbaf8 100644 --- a/README.md +++ b/README.md @@ -99,7 +99,7 @@ Every needle check passed. `--mtp` with the bf16 assistant (`gemma-4-26B-A4B-it- | ~9.8K | 858 / 43.4 | 20.4 GB | 401 / 12.7 | 5.6 GB | | 40.8K | 615 / 36.1 | 21.5 GB | 336 / 12.0 | 5.8 GB | -> ⚠️ **`--stream-experts` crashes on quantized MoE models in releases b769 and b773** (`broadcast_shapes … (N,8,8,D)` on the first request). Reproduced on M6 with Qwen3.6-35B-A3B and Gemma 4 26B-A4B. The mlx-swift-lm upstream sync in #167 broke the SSD path. Earlier versions of this table were measured before that sync and were never re-checked afterwards. Fixed in SharpAI/mlx-swift-lm#69 and #71; the table above was re-measured with those fixes. +> ⚠️ **`--stream-experts` crashes on quantized MoE models in releases b769 and b773** (a `broadcast_shapes` fatal error on the first request, e.g. `(263,8,8,2048) and (263,8,1)` with top-k 8). Reproduced on M6 with Qwen3.6-35B-A3B and Gemma 4 26B-A4B. The mlx-swift-lm upstream sync in #167 broke the SSD path. Earlier versions of this table were measured before that sync and were never re-checked afterwards. Fixed in SharpAI/mlx-swift-lm#69 and #71; the table above was re-measured with those fixes. ### Qwen3.8-27B-4bit (dense) diff --git a/Sources/SwiftLM/Server.swift b/Sources/SwiftLM/Server.swift index 2067d79..cdc1925 100644 --- a/Sources/SwiftLM/Server.swift +++ b/Sources/SwiftLM/Server.swift @@ -277,6 +277,7 @@ final class ProgressTracker { private var lastUpdate: TimeInterval = 0 private var lastBytes: Int64 = 0 private var speedStr = "0.0 MB/s" + private var progress: Progress? init(modelId: String) { self.modelId = modelId @@ -291,6 +292,11 @@ final class ProgressTracker { isDone = true trackingTask?.cancel() if Self.barOpen { + // The Task may be cancelled before its last frame; show where it ended. + if let progress { + print(frame(fraction: progress.fractionCompleted).padding( + toLength: 100, withPad: " ", startingAt: 0), terminator: "") + } print("") fflush(stdout) Self.barOpen = false @@ -320,7 +326,37 @@ final class ProgressTracker { return sumDir(modelHubDir) + sumDir(downloadDir) } + /// One `\r` frame of the download bar at `fraction`. + private func frame(fraction: Double) -> String { + let pct = Int(fraction * 100) + var completedMB = String(format: "%.1f", Double(self.lastBytes) / 1_048_576) + var totalMB = "???" + if fraction > 0.001 { + let extrapolated = (Double(self.lastBytes) / fraction) / 1_048_576.0 + totalMB = String(format: "%.1f", extrapolated) + } else if fraction == 0.0 { + completedMB = "0.0" + } + + let barLength = 20 + let completedBars = min(barLength, Int(fraction * Double(barLength))) + let emptyBars = max(0, barLength - completedBars) + + var bars = "" + if completedBars > 0 { + bars += String(repeating: "=", count: completedBars - 1) + ">" + } + bars += String(repeating: " ", count: emptyBars) + + let pctStr = String(format: "%3d%%", pct) + let spinner = self.spinnerFrames[self.frameIndex] + let speedText = "| Speed: \(self.speedStr)" + + return String(format: "\r[SwiftLM] Download: [%@] %@ %@ (%@ MB / %@ MB) %@", bars, pctStr, spinner, completedMB, totalMB, speedText) + } + func printProgress(_ progress: Progress) { + self.progress = progress if trackingTask == nil { lastUpdate = Date().timeIntervalSince1970 lastBytes = getDownloadedBytes() @@ -329,7 +365,6 @@ final class ProgressTracker { while !self.isDone && !Task.isCancelled { let now = Date().timeIntervalSince1970 let fraction = progress.fractionCompleted - let pct = Int(fraction * 100) let interval = now - self.lastUpdate if interval >= 0.25 { @@ -348,30 +383,7 @@ final class ProgressTracker { self.lastUpdate = now } - var completedMB = String(format: "%.1f", Double(self.lastBytes) / 1_048_576) - var totalMB = "???" - if fraction > 0.001 { - let extrapolated = (Double(self.lastBytes) / fraction) / 1_048_576.0 - totalMB = String(format: "%.1f", extrapolated) - } else if fraction == 0.0 { - completedMB = "0.0" - } - - let barLength = 20 - let completedBars = min(barLength, Int(fraction * Double(barLength))) - let emptyBars = max(0, barLength - completedBars) - - var bars = "" - if completedBars > 0 { - bars += String(repeating: "=", count: completedBars - 1) + ">" - } - bars += String(repeating: " ", count: emptyBars) - - let pctStr = String(format: "%3d%%", pct) - let spinner = self.spinnerFrames[self.frameIndex] - let speedText = "| Speed: \(self.speedStr)" - - let msg = String(format: "\r[SwiftLM] Download: [%@] %@ %@ (%@ MB / %@ MB) %@", bars, pctStr, spinner, completedMB, totalMB, speedText) + let msg = self.frame(fraction: fraction) if self.isDone { break } print(msg.padding(toLength: 100, withPad: " ", startingAt: 0), terminator: "") @@ -439,6 +451,24 @@ func emitEvent(_ payload: [String: Any]) { fflush(stdout) } +/// A `--stream-experts` prefetch finished without every weight shard on disk, +/// for example offline with a partial copy. +struct ModelDownloadIncomplete: LocalizedError { + let modelId: String + let directory: URL + var errorDescription: String? { + "Download of \(modelId) is incomplete (\(directory.path)). Check the network, or delete the directory and retry." + } +} + +/// Whether `directory` has the files `TransformersTokenizerLoader` reads. +func hasTokenizerFiles(in directory: URL) -> Bool { + let fm = FileManager.default + return fm.fileExists(atPath: directory.appendingPathComponent("tokenizer_config.json").path) + && (fm.fileExists(atPath: directory.appendingPathComponent("tokenizer.json").path) + || fm.fileExists(atPath: directory.appendingPathComponent("vocab.json").path)) +} + /// Where `run()` is when it throws, so `classifyExitReason` can produce a more /// precise `reason` than "something failed" and `detail` can name which load /// stalled. Set immediately before each fallible stage begins; read only if @@ -698,9 +728,11 @@ struct MLXServer: AsyncParsableCommand { .appendingPathComponent("MLX", isDirectory: true) .appendingPathComponent("HuggingFace", isDirectory: true)) let localRepo = hub.localRepoLocation(Hub.Repo(id: modelId)) - // Every shard must be present: an interrupted download (or the config.json - // the architecture probe fetches) would otherwise plan with a partial size. - if FileManager.default.fileExists(atPath: localRepo.path), + // Reuse only a copy a finished snapshot left behind. Shards alone are not + // enough: an interrupted download can still be missing tokenizer or template + // files, and the loader won't fill them in when it reads a directory. + let completeMarker = localRepo.appendingPathComponent(".swiftlm-snapshot-complete") + if FileManager.default.fileExists(atPath: completeMarker.path), ModelStorage.validateLocalModelDirectory(localRepo) { modelDirectory = localRepo @@ -710,18 +742,21 @@ struct MLXServer: AsyncParsableCommand { print("[SwiftLM] --stream-experts: downloading \(modelId) before loading...") let prefetchTracker = ProgressTracker(modelId: modelId) defer { prefetchTracker.finish() } - modelDirectory = try await hub.snapshot( + let snapshot = try await hub.snapshot( from: modelId, matching: ["*.safetensors", "*.json", "*.jinja"] ) { progress in prefetchTracker.printProgress(progress) } + // Offline, snapshot returns the repo directory even when it is partial. + guard ModelStorage.validateLocalModelDirectory(snapshot), + hasTokenizerFiles(in: snapshot) + else { + throw ModelDownloadIncomplete(modelId: modelId, directory: snapshot) + } + FileManager.default.createFile(atPath: completeMarker.path, contents: nil) + modelDirectory = snapshot } } - // Streaming is activated for `modelDirectory`, and only a load of that exact - // directory streams. Load from it, or a different lookup could pick another copy. - if self.streamExperts, let dir = modelDirectory { - modelConfig = ModelConfiguration(directory: dir) - } var mainModelProfile: ModelProfile? = nil if self.streamExperts, let dir = modelDirectory { mainModelProfile = ModelProfiler.profile(modelDirectory: dir, modelId: modelId) @@ -739,6 +774,13 @@ struct MLXServer: AsyncParsableCommand { } } + // Streaming is activated for `modelDirectory`, and only a load of that exact + // directory streams. Load from it, or a different lookup could pick another copy. + // Keep the id when tokenizer files are missing, so the loader can fetch them. + if self.streamExperts, let dir = modelDirectory, hasTokenizerFiles(in: dir) { + modelConfig = ModelConfiguration(directory: dir) + } + // Inject streaming flag into config to bypass eval(model) if requested if self.streamExperts { modelConfig.lazyLoad = true From 0f27cbc62b13099b01a4b69eb211475772468f47 Mon Sep 17 00:00:00 2001 From: simba Date: Sun, 27 Sep 2026 16:03:42 -0700 Subject: [PATCH 2/7] docs+test: emitEvent doc placement, SIGTERM mid-stream test Follow-ups from the #197 review: - Move exitAfterShutdownRequest() below emitEvent. It sat between emitEvent's /// block and its declaration, so it took over that doc comment and emitEvent had none. - The SIGPIPE note now names exitAfterShutdownRequest() instead of exit(0). - Test 39 (test-server.sh): SIGTERM while a stream is generating must exit 0, put exiting{reason:"requested"} on its own JSON line, and leave no crash report for that PID. Fails on b782 (exit(), SIGSEGV 139); passes with #197. Co-Authored-By: Claude Opus 5.5 --- Sources/SwiftLM/Server.swift | 32 ++++++++++++------------ tests/test-server.sh | 47 ++++++++++++++++++++++++++++++++++++ 2 files changed, 63 insertions(+), 16 deletions(-) diff --git a/Sources/SwiftLM/Server.swift b/Sources/SwiftLM/Server.swift index cdc1925..d3d32ff 100644 --- a/Sources/SwiftLM/Server.swift +++ b/Sources/SwiftLM/Server.swift @@ -420,21 +420,6 @@ final class ProgressTracker { /// for a protocol-critical signal the daemon is meant to positively rely /// on (not just infer), a dropped emit with zero diagnostic trail would be /// hard to ever notice. Logs to stderr on failure instead. -/// Ends the process after a requested shutdown (SIGTERM/SIGINT) without running -/// C++ static destructors or `atexit` handlers. -/// -/// `exit()` tears those down while an inference thread may still be inside an MLX -/// GPU eval, so a shutdown during generation crashed (SIGSEGV in -/// `CustomKernel::eval_gpu`'s kernel map, or SIGABRT) right after the -/// `exiting{reason:"requested"}` event: the daemon saw a crash report and a -/// non-zero status for a clean stop. Nothing needs those destructors at this -/// point; stdout/stderr are flushed so the exiting event and logs are not lost. -func exitAfterShutdownRequest() -> Never { - fflush(stdout) - fflush(stderr) - Darwin._exit(0) -} - func emitEvent(_ payload: [String: Any]) { guard let data = try? JSONSerialization.data(withJSONObject: payload), let json = String(data: data, encoding: .utf8) @@ -451,6 +436,21 @@ func emitEvent(_ payload: [String: Any]) { fflush(stdout) } +/// Ends the process after a requested shutdown (SIGTERM/SIGINT) without running +/// C++ static destructors or `atexit` handlers. +/// +/// `exit()` tears those down while an inference thread may still be inside an MLX +/// GPU eval, so a shutdown during generation crashed (SIGSEGV in +/// `CustomKernel::eval_gpu`'s kernel map, or SIGABRT) right after the +/// `exiting{reason:"requested"}` event: the daemon saw a crash report and a +/// non-zero status for a clean stop. Nothing needs those destructors at this +/// point; stdout/stderr are flushed so the exiting event and logs are not lost. +func exitAfterShutdownRequest() -> Never { + fflush(stdout) + fflush(stderr) + Darwin._exit(0) +} + /// A `--stream-experts` prefetch finished without every weight shard on disk, /// for example offline with a partial copy. struct ModelDownloadIncomplete: LocalizedError { @@ -1647,7 +1647,7 @@ struct MLXServer: AsyncParsableCommand { // is already gone by the time a shutdown signal arrives (e.g. the // daemon itself already crashed), the exiting-event print()/fflush // below can raise SIGPIPE — whose default disposition kills this - // process via signal instead of reaching Darwin.exit(0), producing + // process via signal instead of reaching exitAfterShutdownRequest(), producing // exactly the ambiguous "was this a crash?" signature this feature // exists to eliminate. Ignore SIGPIPE so a closed pipe surfaces as // an ordinary EPIPE write error instead. diff --git a/tests/test-server.sh b/tests/test-server.sh index 95b0939..e07cb46 100755 --- a/tests/test-server.sh +++ b/tests/test-server.sh @@ -35,6 +35,10 @@ cleanup() { kill -9 "$SERVER_PID" 2>/dev/null || true wait "$SERVER_PID" 2>/dev/null || true fi + if [ -n "${TERM_SERVER_PID:-}" ]; then + kill -9 "$TERM_SERVER_PID" 2>/dev/null || true + wait "$TERM_SERVER_PID" 2>/dev/null || true + fi if [ -n "${CORS_SERVER_PID:-}" ]; then log "Stopping CORS server (PID $CORS_SERVER_PID)" kill -9 "$CORS_SERVER_PID" 2>/dev/null || true @@ -1203,6 +1207,49 @@ fi rm -f /tmp/mlx_ttfb_prompt.txt /tmp/mlx_ttfb_body.json /tmp/mlx_ttfb_stream.txt +# ── Test 39: SIGTERM during generation exits cleanly ───────────────── +# A shutdown request mid-stream must exit 0 with exiting{reason:"requested"} on its +# own line, and must not crash (exit() used to race an in-flight GPU eval, #197). +log "Test 39: SIGTERM during a streaming generation" +TERM_PORT=$((PORT + 3)) +TERM_LOG=$(mktemp) +TERM_STREAM=$(mktemp) +TERM_MARK=$(mktemp) # only crash reports written after this can be ours +"$BINARY" --model "$MODEL" --port "$TERM_PORT" --host "$HOST" > "$TERM_LOG" 2>&1 & +TERM_SERVER_PID=$! +TERM_PID=$TERM_SERVER_PID +for i in $(seq 1 120); do + curl -sf "http://${HOST}:${TERM_PORT}/health" >/dev/null 2>&1 && break + sleep 1 +done +curl -sN -X POST "http://${HOST}:${TERM_PORT}/v1/chat/completions" \ + -H "Content-Type: application/json" \ + -d "{\"model\":\"$MODEL\",\"stream\":true,\"max_tokens\":2048,\"messages\":[{\"role\":\"user\",\"content\":\"Write a very long story about a lighthouse keeper.\"}]}" \ + > "$TERM_STREAM" 2>/dev/null & +TERM_CURL_PID=$! +for i in $(seq 1 200); do + grep -q '"content"' "$TERM_STREAM" 2>/dev/null && break + sleep 0.1 +done +kill -TERM "$TERM_SERVER_PID" +TERM_STATUS=0 +wait "$TERM_SERVER_PID" || TERM_STATUS=$? +unset TERM_SERVER_PID +kill "$TERM_CURL_PID" 2>/dev/null || true +wait "$TERM_CURL_PID" 2>/dev/null || true + +TERM_REASON=$(grep '^{' "$TERM_LOG" | jq -r 'select(.event == "exiting") | .reason' 2>/dev/null | tail -1) +sleep 3 # ReportCrash writes the .ips a moment after the process dies +TERM_CRASHES=$(find "$HOME/Library/Logs/DiagnosticReports" -name 'SwiftLM*' -newer "$TERM_MARK" \ + -exec grep -l "\"pid\" : ${TERM_PID}," {} + 2>/dev/null | wc -l | tr -d ' ') +if [ "$TERM_STATUS" -eq 0 ] && [ "$TERM_REASON" = "requested" ] && [ "$TERM_CRASHES" -eq 0 ]; then + pass "SIGTERM mid-stream: exit 0, exiting{requested} on its own line, no crash report" +else + fail "SIGTERM mid-stream: status=$TERM_STATUS reason='${TERM_REASON}' crash_reports=$TERM_CRASHES" + tail -5 "$TERM_LOG" +fi +rm -f "$TERM_LOG" "$TERM_STREAM" "$TERM_MARK" + # ── Results ────────────────────────────────────────────────────────── echo "" log "═══════════════════════════════════════" From 3284dc5088193ee5b2d6b993aad15ac549c517eb Mon Sep 17 00:00:00 2001 From: simba Date: Sun, 27 Sep 2026 18:21:24 -0700 Subject: [PATCH 3/7] fix(ssd): review fixes for #198 - A resolved model directory without tokenizer.json is treated as not found under --stream-experts, so the prefetch builds a complete copy and streaming targets the directory that is loaded (was a re-download plus a non-streaming load). - Check tokenizer.json, which swift-transformers requires; tokenizer_config.json is optional and vocab.json is no substitute. A repo without it fails with ModelMissingTokenizer instead of 'download incomplete'. - Write the completion marker only when the online file list succeeds and every listed file exists, so an offline partial snapshot never pins itself. - If the Hub is unreachable but a loadable local copy exists (for example one from before the marker), warn and use it instead of failing. - Finish the prefetch bar as soon as snapshot returns, before validation output. - ModelDownloadIncomplete/ModelMissingTokenizer are CustomStringConvertible, so the exiting event's detail is the message, not a struct dump. - The shutdown notice and the exiting JSON go out in one print, so a generation token can't land between them and push the JSON off its line. - Test 39: nothing in it can abort the script under set -euo pipefail; it fails cleanly if the server never gets ready, bounds the shutdown wait at 30 s, and drops the .ips check (reports arrive too late; a crash already shows as a nonzero status). Co-Authored-By: Claude Opus 5.5 --- Sources/SwiftLM/Server.swift | 103 +++++++++++++++++++++++------------ tests/test-server.sh | 87 ++++++++++++++++++----------- 2 files changed, 125 insertions(+), 65 deletions(-) diff --git a/Sources/SwiftLM/Server.swift b/Sources/SwiftLM/Server.swift index d3d32ff..d97f314 100644 --- a/Sources/SwiftLM/Server.swift +++ b/Sources/SwiftLM/Server.swift @@ -420,7 +420,9 @@ final class ProgressTracker { /// for a protocol-critical signal the daemon is meant to positively rely /// on (not just infer), a dropped emit with zero diagnostic trail would be /// hard to ever notice. Logs to stderr on failure instead. -func emitEvent(_ payload: [String: Any]) { +/// `preface` lines are written in the same `print` as the JSON, so no generation +/// token can land between them and push the JSON off the start of its line. +func emitEvent(_ payload: [String: Any], preface: String? = nil) { guard let data = try? JSONSerialization.data(withJSONObject: payload), let json = String(data: data, encoding: .utf8) else { @@ -428,11 +430,13 @@ func emitEvent(_ payload: [String: Any]) { Data("[SwiftLM] failed to encode event for stdout: \(payload)\n".utf8)) return } + var out = "" if ProgressTracker.barOpen { - print("") + out += "\n" ProgressTracker.barOpen = false } - print(json) + if let preface { out += preface + "\n" } + print(out + json) fflush(stdout) } @@ -453,20 +457,32 @@ func exitAfterShutdownRequest() -> Never { /// A `--stream-experts` prefetch finished without every weight shard on disk, /// for example offline with a partial copy. -struct ModelDownloadIncomplete: LocalizedError { +struct ModelDownloadIncomplete: Error, CustomStringConvertible { let modelId: String let directory: URL - var errorDescription: String? { + var description: String { "Download of \(modelId) is incomplete (\(directory.path)). Check the network, or delete the directory and retry." } } -/// Whether `directory` has the files `TransformersTokenizerLoader` reads. -func hasTokenizerFiles(in directory: URL) -> Bool { - let fm = FileManager.default - return fm.fileExists(atPath: directory.appendingPathComponent("tokenizer_config.json").path) - && (fm.fileExists(atPath: directory.appendingPathComponent("tokenizer.json").path) - || fm.fileExists(atPath: directory.appendingPathComponent("vocab.json").path)) +/// The repository has no `tokenizer.json`, which the tokenizer loader requires. +struct ModelMissingTokenizer: Error, CustomStringConvertible { + let modelId: String + var description: String { + "\(modelId) has no tokenizer.json, which SwiftLM needs to load it." + } +} + +/// Whether `directory` has `tokenizer.json`. swift-transformers requires it; +/// `tokenizer_config.json` is optional, and `vocab.json` is no substitute. +func hasTokenizerJSON(in directory: URL) -> Bool { + FileManager.default.fileExists(atPath: directory.appendingPathComponent("tokenizer.json").path) +} + +/// Whether a local copy can be streamed and loaded by directory: every shard, plus +/// the tokenizer the loader won't fetch when it reads a directory. +func isLoadableModelDirectory(_ directory: URL) -> Bool { + ModelStorage.validateLocalModelDirectory(directory) && hasTokenizerJSON(in: directory) } /// Where `run()` is when it throws, so `classifyExitReason` can produce a more @@ -712,10 +728,9 @@ struct MLXServer: AsyncParsableCommand { ModelStorage.validatedContentDirectory(for: modelId) ?? resolveModelDirectory(modelId: modelId) // resolveModelDirectory doesn't check the weights are there; streaming must not be - // activated for an empty or partial snapshot the loader won't read. - if self.streamExperts, let dir = modelDirectory, - !ModelStorage.validateLocalModelDirectory(dir) - { + // activated for an empty or partial snapshot, or one without a tokenizer, since + // the streamed model loads from this exact directory and nothing fills it in. + if self.streamExperts, let dir = modelDirectory, !isLoadableModelDirectory(dir) { modelDirectory = nil } if self.streamExperts, !self.info, modelDirectory == nil, @@ -732,8 +747,9 @@ struct MLXServer: AsyncParsableCommand { // enough: an interrupted download can still be missing tokenizer or template // files, and the loader won't fill them in when it reads a directory. let completeMarker = localRepo.appendingPathComponent(".swiftlm-snapshot-complete") + let patterns = ["*.safetensors", "*.json", "*.jinja"] if FileManager.default.fileExists(atPath: completeMarker.path), - ModelStorage.validateLocalModelDirectory(localRepo) + isLoadableModelDirectory(localRepo) { modelDirectory = localRepo } else { @@ -742,19 +758,37 @@ struct MLXServer: AsyncParsableCommand { print("[SwiftLM] --stream-experts: downloading \(modelId) before loading...") let prefetchTracker = ProgressTracker(modelId: modelId) defer { prefetchTracker.finish() } - let snapshot = try await hub.snapshot( - from: modelId, matching: ["*.safetensors", "*.json", "*.jinja"] - ) { progress in - prefetchTracker.printProgress(progress) - } - // Offline, snapshot returns the repo directory even when it is partial. - guard ModelStorage.validateLocalModelDirectory(snapshot), - hasTokenizerFiles(in: snapshot) - else { - throw ModelDownloadIncomplete(modelId: modelId, directory: snapshot) + do { + let snapshot = try await hub.snapshot(from: modelId, matching: patterns) { + progress in prefetchTracker.printProgress(progress) + } + prefetchTracker.finish() + // Offline, snapshot returns the repo directory even when it is partial. + guard ModelStorage.validateLocalModelDirectory(snapshot) else { + throw ModelDownloadIncomplete(modelId: modelId, directory: snapshot) + } + guard hasTokenizerJSON(in: snapshot) else { + throw ModelMissingTokenizer(modelId: modelId) + } + // Mark complete only against the online file list, so an offline + // partial copy is healed by the next online start. + if let names = try? await hub.getFilenames(from: modelId, matching: patterns), + !names.isEmpty, + names.allSatisfy({ + FileManager.default.fileExists( + atPath: snapshot.appendingPathComponent($0).path) + }) + { + FileManager.default.createFile(atPath: completeMarker.path, contents: nil) + } + modelDirectory = snapshot + } catch where isLoadableModelDirectory(localRepo) { + // Hub unreachable, but a loadable copy is already here (for example + // one from before the marker existed): use it rather than fail. + prefetchTracker.finish() + print("[SwiftLM] ⚠️ Could not verify \(modelId) with the Hub (\(error)); using the local copy.") + modelDirectory = localRepo } - FileManager.default.createFile(atPath: completeMarker.path, contents: nil) - modelDirectory = snapshot } } var mainModelProfile: ModelProfile? = nil @@ -776,8 +810,7 @@ struct MLXServer: AsyncParsableCommand { // Streaming is activated for `modelDirectory`, and only a load of that exact // directory streams. Load from it, or a different lookup could pick another copy. - // Keep the id when tokenizer files are missing, so the loader can fetch them. - if self.streamExperts, let dir = modelDirectory, hasTokenizerFiles(in: dir) { + if self.streamExperts, let dir = modelDirectory { modelConfig = ModelConfiguration(directory: dir) } @@ -1654,13 +1687,15 @@ struct MLXServer: AsyncParsableCommand { signal(SIGPIPE, SIG_IGN) shutdownSource.setEventHandler { - print("\n[SwiftLM] Received SIGTERM, shutting down gracefully...") - emitEvent(["event": "exiting", "reason": "requested"]) + emitEvent( + ["event": "exiting", "reason": "requested"], + preface: "\n[SwiftLM] Received SIGTERM, shutting down gracefully...") exitAfterShutdownRequest() } interruptSource.setEventHandler { - print("\n[SwiftLM] Received SIGINT, shutting down gracefully...") - emitEvent(["event": "exiting", "reason": "requested"]) + emitEvent( + ["event": "exiting", "reason": "requested"], + preface: "\n[SwiftLM] Received SIGINT, shutting down gracefully...") exitAfterShutdownRequest() } shutdownSource.resume() diff --git a/tests/test-server.sh b/tests/test-server.sh index e07cb46..8dbde21 100755 --- a/tests/test-server.sh +++ b/tests/test-server.sh @@ -1209,46 +1209,71 @@ rm -f /tmp/mlx_ttfb_prompt.txt /tmp/mlx_ttfb_body.json /tmp/mlx_ttfb_stream.txt # ── Test 39: SIGTERM during generation exits cleanly ───────────────── # A shutdown request mid-stream must exit 0 with exiting{reason:"requested"} on its -# own line, and must not crash (exit() used to race an in-flight GPU eval, #197). +# own line (exit() used to race an in-flight GPU eval and crash, #197). A crash +# shows up as a nonzero status; crash reports are written too late to check here. log "Test 39: SIGTERM during a streaming generation" TERM_PORT=$((PORT + 3)) TERM_LOG=$(mktemp) TERM_STREAM=$(mktemp) -TERM_MARK=$(mktemp) # only crash reports written after this can be ours "$BINARY" --model "$MODEL" --port "$TERM_PORT" --host "$HOST" > "$TERM_LOG" 2>&1 & TERM_SERVER_PID=$! -TERM_PID=$TERM_SERVER_PID +TERM_READY=false for i in $(seq 1 120); do - curl -sf "http://${HOST}:${TERM_PORT}/health" >/dev/null 2>&1 && break + if curl -sf "http://${HOST}:${TERM_PORT}/health" >/dev/null 2>&1; then + TERM_READY=true + break + fi + kill -0 "$TERM_SERVER_PID" 2>/dev/null || break sleep 1 done -curl -sN -X POST "http://${HOST}:${TERM_PORT}/v1/chat/completions" \ - -H "Content-Type: application/json" \ - -d "{\"model\":\"$MODEL\",\"stream\":true,\"max_tokens\":2048,\"messages\":[{\"role\":\"user\",\"content\":\"Write a very long story about a lighthouse keeper.\"}]}" \ - > "$TERM_STREAM" 2>/dev/null & -TERM_CURL_PID=$! -for i in $(seq 1 200); do - grep -q '"content"' "$TERM_STREAM" 2>/dev/null && break - sleep 0.1 -done -kill -TERM "$TERM_SERVER_PID" -TERM_STATUS=0 -wait "$TERM_SERVER_PID" || TERM_STATUS=$? -unset TERM_SERVER_PID -kill "$TERM_CURL_PID" 2>/dev/null || true -wait "$TERM_CURL_PID" 2>/dev/null || true - -TERM_REASON=$(grep '^{' "$TERM_LOG" | jq -r 'select(.event == "exiting") | .reason' 2>/dev/null | tail -1) -sleep 3 # ReportCrash writes the .ips a moment after the process dies -TERM_CRASHES=$(find "$HOME/Library/Logs/DiagnosticReports" -name 'SwiftLM*' -newer "$TERM_MARK" \ - -exec grep -l "\"pid\" : ${TERM_PID}," {} + 2>/dev/null | wc -l | tr -d ' ') -if [ "$TERM_STATUS" -eq 0 ] && [ "$TERM_REASON" = "requested" ] && [ "$TERM_CRASHES" -eq 0 ]; then - pass "SIGTERM mid-stream: exit 0, exiting{requested} on its own line, no crash report" -else - fail "SIGTERM mid-stream: status=$TERM_STATUS reason='${TERM_REASON}' crash_reports=$TERM_CRASHES" - tail -5 "$TERM_LOG" -fi -rm -f "$TERM_LOG" "$TERM_STREAM" "$TERM_MARK" + +if [ "$TERM_READY" != true ]; then + fail "SIGTERM mid-stream: server did not become ready" + tail -5 "$TERM_LOG" || true + kill -9 "$TERM_SERVER_PID" 2>/dev/null || true + wait "$TERM_SERVER_PID" 2>/dev/null || true + unset TERM_SERVER_PID +else + curl -sN -X POST "http://${HOST}:${TERM_PORT}/v1/chat/completions" \ + -H "Content-Type: application/json" \ + -d "{\"model\":\"$MODEL\",\"stream\":true,\"max_tokens\":2048,\"messages\":[{\"role\":\"user\",\"content\":\"Write a very long story about a lighthouse keeper.\"}]}" \ + > "$TERM_STREAM" 2>/dev/null & + TERM_CURL_PID=$! + for i in $(seq 1 200); do + grep -q '"content"' "$TERM_STREAM" 2>/dev/null && break + sleep 0.1 + done + kill -TERM "$TERM_SERVER_PID" 2>/dev/null || true + + # Bounded wait: a shutdown hang must fail here, not at the job timeout. + TERM_HUNG=false + for i in $(seq 1 300); do + kill -0 "$TERM_SERVER_PID" 2>/dev/null || break + sleep 0.1 + done + if kill -0 "$TERM_SERVER_PID" 2>/dev/null; then + TERM_HUNG=true + kill -9 "$TERM_SERVER_PID" 2>/dev/null || true + fi + TERM_STATUS=0 + wait "$TERM_SERVER_PID" 2>/dev/null || TERM_STATUS=$? + unset TERM_SERVER_PID + kill "$TERM_CURL_PID" 2>/dev/null || true + wait "$TERM_CURL_PID" 2>/dev/null || true + + TERM_REASON=$({ grep '^{' "$TERM_LOG" || true; } \ + | jq -r 'select(.event == "exiting") | .reason' 2>/dev/null | tail -1 || true) + if [ "$TERM_HUNG" = true ]; then + fail "SIGTERM mid-stream: server still running 30s after SIGTERM" + tail -5 "$TERM_LOG" || true + elif [ "$TERM_STATUS" -eq 0 ] && [ "$TERM_REASON" = "requested" ]; then + pass "SIGTERM mid-stream: exit 0, exiting{requested} on its own line" + else + fail "SIGTERM mid-stream: status=$TERM_STATUS reason='${TERM_REASON}'" + tail -5 "$TERM_LOG" || true + fi +fi +rm -f "$TERM_LOG" "$TERM_STREAM" # ── Results ────────────────────────────────────────────────────────── echo "" From 4547bce2d1fc590795e97af608233f73810c954d Mon Sep 17 00:00:00 2001 From: simba Date: Sun, 27 Sep 2026 20:22:54 -0700 Subject: [PATCH 4/7] fix(ssd): decide local-copy completeness from the Hub's file list Review round 2 for #198. The prefetch guessed completeness from the shard layout, which rejected valid repos and trusted partial ones. It now asks the Hub which files the repo has (StreamingDirectory.swift): - Online, a local copy is used only if it has every listed file: the ~/.cache copy first, then the loader's Application Support copy. Otherwise the model is downloaded and must then have them all. This accepts any weight layout the loader does (weights.NN.safetensors, no index) and stops an interrupted download in ~/.cache from being reused forever. - A repo without tokenizer.json fails before downloading, with its own message. - A Hub error (unknown or gated id) fails the start with a clear message; only network errors and timeouts count as offline. - Offline, only a copy known to be complete is used: the marker written after a listing check, or a quietly validated copy with tokenizer.json. - A dense model found locally isn't downloaded in full first; streaming turns off for it and the loader completes it. - --info and local paths are left as they were. - Non-streaming starts also check a validated local copy against the listing (best effort) and load through the Hub when files are missing. - Test 39 parses each log line as JSON on its own, keeps the log on failure, and cleanup() removes its files and curl. Co-Authored-By: Claude Opus 5.5 --- Sources/MLXInferenceCore/ModelStorage.swift | 5 + Sources/SwiftLM/Server.swift | 124 ++++------------ Sources/SwiftLM/StreamingDirectory.swift | 157 ++++++++++++++++++++ tests/test-server.sh | 21 ++- 4 files changed, 208 insertions(+), 99 deletions(-) create mode 100644 Sources/SwiftLM/StreamingDirectory.swift diff --git a/Sources/MLXInferenceCore/ModelStorage.swift b/Sources/MLXInferenceCore/ModelStorage.swift index 9c39c85..da7429e 100644 --- a/Sources/MLXInferenceCore/ModelStorage.swift +++ b/Sources/MLXInferenceCore/ModelStorage.swift @@ -280,6 +280,11 @@ public enum ModelStorage { validateModelFiles(in: directory, logFailures: true) } + /// `validateLocalModelDirectory` without logging, for probing candidate copies. + public static func validateLocalModelDirectory(_ directory: URL, logFailures: Bool) -> Bool { + validateModelFiles(in: directory, logFailures: logFailures) + } + /// Read the raw config.json dictionary for a downloaded model. /// Verifies that all required safetensors files are present in the snapshot directory. /// This prevents the engine from entering `.ready` state if a download was interrupted or corrupted. diff --git a/Sources/SwiftLM/Server.swift b/Sources/SwiftLM/Server.swift index d97f314..3c097d8 100644 --- a/Sources/SwiftLM/Server.swift +++ b/Sources/SwiftLM/Server.swift @@ -455,36 +455,6 @@ func exitAfterShutdownRequest() -> Never { Darwin._exit(0) } -/// A `--stream-experts` prefetch finished without every weight shard on disk, -/// for example offline with a partial copy. -struct ModelDownloadIncomplete: Error, CustomStringConvertible { - let modelId: String - let directory: URL - var description: String { - "Download of \(modelId) is incomplete (\(directory.path)). Check the network, or delete the directory and retry." - } -} - -/// The repository has no `tokenizer.json`, which the tokenizer loader requires. -struct ModelMissingTokenizer: Error, CustomStringConvertible { - let modelId: String - var description: String { - "\(modelId) has no tokenizer.json, which SwiftLM needs to load it." - } -} - -/// Whether `directory` has `tokenizer.json`. swift-transformers requires it; -/// `tokenizer_config.json` is optional, and `vocab.json` is no substitute. -func hasTokenizerJSON(in directory: URL) -> Bool { - FileManager.default.fileExists(atPath: directory.appendingPathComponent("tokenizer.json").path) -} - -/// Whether a local copy can be streamed and loaded by directory: every shard, plus -/// the tokenizer the loader won't fetch when it reads a directory. -func isLoadableModelDirectory(_ directory: URL) -> Bool { - ModelStorage.validateLocalModelDirectory(directory) && hasTokenizerJSON(in: directory) -} - /// Where `run()` is when it throws, so `classifyExitReason` can produce a more /// precise `reason` than "something failed" and `detail` can name which load /// stalled. Set immediately before each fallible stage begins; read only if @@ -697,11 +667,37 @@ struct MLXServer: AsyncParsableCommand { let modelId = model // ── Load model ── + // Same root the loader's HubApi uses further down. + let cliHub = HubApi( + downloadBase: URL.applicationSupportDirectory + .appendingPathComponent("MLX", isDirectory: true) + .appendingPathComponent("HuggingFace", isDirectory: true)) + let isHubId = !ModelStorage.isLocalDirectoryPath(modelId) + && !FileManager.default.fileExists(atPath: modelId) + let validatedLocal = isHubId ? ModelStorage.validatedContentDirectory(for: modelId) : nil + // A copy whose download stopped after the shards still validates, so check local + // copies against the Hub's file list. --stream-experts needs the list anyway, and + // a Hub error there (unknown or gated id) fails the start. + var hubListing: HubListing? = nil + if isHubId, !self.info, self.streamExperts || validatedLocal != nil { + phase = .architectureProbe + hubListing = self.streamExperts + ? try await fetchHubListing(cliHub, modelId: modelId) + : try? await fetchHubListing(cliHub, modelId: modelId) + } + var validatedLocalMissing: [String] = [] + if let dir = validatedLocal, case .files(let files)? = hubListing { + validatedLocalMissing = missingFiles(files, in: dir) + } + var modelConfig: ModelConfiguration if ModelStorage.isLocalDirectoryPath(modelId) { print("[SwiftLM] Loading from local directory: \(modelId)") modelConfig = ModelConfiguration(directory: URL(filePath: modelId)) - } else if let localDirectory = ModelStorage.validatedContentDirectory(for: modelId) { + } else if let dir = validatedLocal, !validatedLocalMissing.isEmpty { + print("[SwiftLM] \(dir.path) is missing \(validatedLocalMissing.count) file(s) (e.g. \(validatedLocalMissing[0])); loading through the Hub to complete it.") + modelConfig = ModelConfiguration(id: modelId) + } else if let localDirectory = validatedLocal { // Any validated copy in the shared HF cache, in any supported layout. Note // this deliberately does NOT use localLoadDirectory: that skips the // materialized `models//` layout on the grounds that HubApi @@ -727,69 +723,9 @@ struct MLXServer: AsyncParsableCommand { var modelDirectory = ModelStorage.validatedContentDirectory(for: modelId) ?? resolveModelDirectory(modelId: modelId) - // resolveModelDirectory doesn't check the weights are there; streaming must not be - // activated for an empty or partial snapshot, or one without a tokenizer, since - // the streamed model loads from this exact directory and nothing fills it in. - if self.streamExperts, let dir = modelDirectory, !isLoadableModelDirectory(dir) { - modelDirectory = nil - } - if self.streamExperts, !self.info, modelDirectory == nil, - !FileManager.default.fileExists(atPath: modelId) - { - // Streaming must be activated for the directory the loader reads, so resolve it - // before loading. Same hub root as the loader below, which reuses these files. - let hub = HubApi( - downloadBase: URL.applicationSupportDirectory - .appendingPathComponent("MLX", isDirectory: true) - .appendingPathComponent("HuggingFace", isDirectory: true)) - let localRepo = hub.localRepoLocation(Hub.Repo(id: modelId)) - // Reuse only a copy a finished snapshot left behind. Shards alone are not - // enough: an interrupted download can still be missing tokenizer or template - // files, and the loader won't fill them in when it reads a directory. - let completeMarker = localRepo.appendingPathComponent(".swiftlm-snapshot-complete") - let patterns = ["*.safetensors", "*.json", "*.jinja"] - if FileManager.default.fileExists(atPath: completeMarker.path), - isLoadableModelDirectory(localRepo) - { - modelDirectory = localRepo - } else { - // First run. A failed download is a model problem, not a binary one. - phase = .architectureProbe - print("[SwiftLM] --stream-experts: downloading \(modelId) before loading...") - let prefetchTracker = ProgressTracker(modelId: modelId) - defer { prefetchTracker.finish() } - do { - let snapshot = try await hub.snapshot(from: modelId, matching: patterns) { - progress in prefetchTracker.printProgress(progress) - } - prefetchTracker.finish() - // Offline, snapshot returns the repo directory even when it is partial. - guard ModelStorage.validateLocalModelDirectory(snapshot) else { - throw ModelDownloadIncomplete(modelId: modelId, directory: snapshot) - } - guard hasTokenizerJSON(in: snapshot) else { - throw ModelMissingTokenizer(modelId: modelId) - } - // Mark complete only against the online file list, so an offline - // partial copy is healed by the next online start. - if let names = try? await hub.getFilenames(from: modelId, matching: patterns), - !names.isEmpty, - names.allSatisfy({ - FileManager.default.fileExists( - atPath: snapshot.appendingPathComponent($0).path) - }) - { - FileManager.default.createFile(atPath: completeMarker.path, contents: nil) - } - modelDirectory = snapshot - } catch where isLoadableModelDirectory(localRepo) { - // Hub unreachable, but a loadable copy is already here (for example - // one from before the marker existed): use it rather than fail. - prefetchTracker.finish() - print("[SwiftLM] ⚠️ Could not verify \(modelId) with the Hub (\(error)); using the local copy.") - modelDirectory = localRepo - } - } + if self.streamExperts, !self.info, isHubId, let listing = hubListing { + modelDirectory = try await resolveStreamingDirectory( + modelId: modelId, candidate: modelDirectory, hub: cliHub, listing: listing) } var mainModelProfile: ModelProfile? = nil if self.streamExperts, let dir = modelDirectory { diff --git a/Sources/SwiftLM/StreamingDirectory.swift b/Sources/SwiftLM/StreamingDirectory.swift new file mode 100644 index 0000000..63ee5a3 --- /dev/null +++ b/Sources/SwiftLM/StreamingDirectory.swift @@ -0,0 +1,157 @@ +// StreamingDirectory.swift — where a Hub model is loaded from, checked against the +// Hub's own file list. +// +// `--stream-experts` streams from the one directory it was activated for, and a +// directory load doesn't fill in missing files. So the directory must be complete, +// and "complete" is decided by what the Hub lists, not by guessing a shard layout +// (repos ship `weights.NN.safetensors`, no index, or an index naming extra shards). + +import Foundation +import Hub +import MLXInferenceCore + +/// Files SwiftLM downloads for a model: weights, configs, tokenizer and templates. +let modelDownloadPatterns = ["*.safetensors", "*.json", "*.jinja"] + +/// A `--stream-experts` prefetch finished without every listed file on disk. +struct ModelDownloadIncomplete: Error, CustomStringConvertible { + let modelId: String + let directory: URL + var description: String { + "Download of \(modelId) is incomplete (\(directory.path)). Check the network, or delete the directory and retry." + } +} + +/// The repository has no `tokenizer.json`, which the tokenizer loader requires. +struct ModelMissingTokenizer: Error, CustomStringConvertible { + let modelId: String + var description: String { + "\(modelId) has no tokenizer.json, which SwiftLM needs to load it." + } +} + +/// The Hub couldn't be reached and there is no local copy known to be complete. +struct ModelUnavailableOffline: Error, CustomStringConvertible { + let modelId: String + let underlying: Error + var description: String { + "Could not reach the Hub to fetch \(modelId) (\(underlying)), and no complete local copy was found." + } +} + +/// The Hub answered, but has no repository `modelId` this user can read. +struct ModelNotOnHub: Error, CustomStringConvertible { + let modelId: String + let underlying: Error + var description: String { + "The Hub has no accessible repository \(modelId) (\(underlying)). Check the id, or log in for a gated model." + } +} + +struct HubListingTimeout: Error, CustomStringConvertible { + var description: String { "the Hub did not answer in time" } +} + +/// The Hub's file list for a repository, or why there isn't one. +enum HubListing { + case files([String]) + /// A network error or timeout; the Hub gave no answer. + case unreachable(Error) +} + +/// Asks the Hub which files `modelId` has. Throws only when the Hub answered with +/// an error (unknown or gated repository); connectivity problems are `.unreachable`. +func fetchHubListing(_ hub: HubApi, modelId: String) async throws -> HubListing { + do { + let files = try await withThrowingTaskGroup(of: [String].self) { group in + group.addTask { try await hub.getFilenames(from: modelId, matching: modelDownloadPatterns) } + group.addTask { + try await Task.sleep(nanoseconds: 15_000_000_000) + throw HubListingTimeout() + } + defer { group.cancelAll() } + return try await group.next()! + } + return .files(files) + } catch let error as Hub.HubClientError { + switch error { + case .httpStatusCode, .authorizationRequired, .resourceNotFound, .fileNotFound: + throw ModelNotOnHub(modelId: modelId, underlying: error) + default: + return .unreachable(error) + } + } catch { + return .unreachable(error) + } +} + +/// Listed files that are absent from `directory`. +func missingFiles(_ files: [String], in directory: URL) -> [String] { + files.filter { !FileManager.default.fileExists(atPath: directory.appendingPathComponent($0).path) } +} + +/// Whether `directory` has `tokenizer.json`. swift-transformers requires it; +/// `tokenizer_config.json` is optional, and `vocab.json` is no substitute. +func hasTokenizerJSON(in directory: URL) -> Bool { + FileManager.default.fileExists(atPath: directory.appendingPathComponent("tokenizer.json").path) +} + +/// The directory `--stream-experts` loads and streams `modelId` from. +/// +/// Online, a local copy (`candidate`, then the loader's Application Support copy) +/// is used only if it has every listed file; otherwise the model is downloaded and +/// must then have them all. Offline, only a copy known to be complete is used. +func resolveStreamingDirectory( + modelId: String, candidate: URL?, hub: HubApi, listing: HubListing +) async throws -> URL { + let fm = FileManager.default + let localRepo = hub.localRepoLocation(Hub.Repo(id: modelId)) + // Records that this copy once matched the Hub listing, for offline starts. + let marker = localRepo.appendingPathComponent(".swiftlm-snapshot-complete") + let copies = [candidate, localRepo].compactMap { $0 }.filter { fm.fileExists(atPath: $0.path) } + + switch listing { + case .files(let files): + guard files.contains("tokenizer.json") else { throw ModelMissingTokenizer(modelId: modelId) } + for dir in copies { + let missing = missingFiles(files, in: dir) + if missing.isEmpty { + if dir == localRepo { fm.createFile(atPath: marker.path, contents: nil) } + return dir + } + print("[SwiftLM] \(dir.path) is missing \(missing.count) of \(files.count) files (e.g. \(missing[0])).") + } + // A dense model won't stream: leave completing it to the loader instead of + // downloading everything here first. + if let dir = copies.first(where: { + fm.fileExists(atPath: $0.appendingPathComponent("config.json").path) + }), let profile = ModelProfiler.profile(modelDirectory: dir, modelId: modelId), !profile.isMoE { + return dir + } + + print("[SwiftLM] --stream-experts: downloading \(modelId) before loading...") + let tracker = ProgressTracker(modelId: modelId) + defer { tracker.finish() } + let snapshot = try await hub.snapshot(from: modelId, matching: modelDownloadPatterns) { + tracker.printProgress($0) + } + tracker.finish() + guard missingFiles(files, in: snapshot).isEmpty else { + throw ModelDownloadIncomplete(modelId: modelId, directory: snapshot) + } + fm.createFile(atPath: marker.path, contents: nil) + return snapshot + + case .unreachable(let error): + let verified = fm.fileExists(atPath: marker.path) && hasTokenizerJSON(in: localRepo) + ? [localRepo] : [] + let plausible = copies.filter { + hasTokenizerJSON(in: $0) && ModelStorage.validateLocalModelDirectory($0, logFailures: false) + } + guard let dir = (verified + plausible).first else { + throw ModelUnavailableOffline(modelId: modelId, underlying: error) + } + print("[SwiftLM] ⚠️ Could not reach the Hub to check \(modelId) (\(error)); using \(dir.path).") + return dir + } +} diff --git a/tests/test-server.sh b/tests/test-server.sh index 8dbde21..36067b7 100755 --- a/tests/test-server.sh +++ b/tests/test-server.sh @@ -39,6 +39,10 @@ cleanup() { kill -9 "$TERM_SERVER_PID" 2>/dev/null || true wait "$TERM_SERVER_PID" 2>/dev/null || true fi + if [ -n "${TERM_CURL_PID:-}" ]; then + kill "$TERM_CURL_PID" 2>/dev/null || true + fi + rm -f "${TERM_LOG:-}" "${TERM_STREAM:-}" if [ -n "${CORS_SERVER_PID:-}" ]; then log "Stopping CORS server (PID $CORS_SERVER_PID)" kill -9 "$CORS_SERVER_PID" 2>/dev/null || true @@ -1229,7 +1233,8 @@ done if [ "$TERM_READY" != true ]; then fail "SIGTERM mid-stream: server did not become ready" - tail -5 "$TERM_LOG" || true + tail -40 "$TERM_LOG" || true + cp "$TERM_LOG" /tmp/SwiftLM-test-sigterm.log 2>/dev/null || true kill -9 "$TERM_SERVER_PID" 2>/dev/null || true wait "$TERM_SERVER_PID" 2>/dev/null || true unset TERM_SERVER_PID @@ -1260,20 +1265,26 @@ else unset TERM_SERVER_PID kill "$TERM_CURL_PID" 2>/dev/null || true wait "$TERM_CURL_PID" 2>/dev/null || true + unset TERM_CURL_PID - TERM_REASON=$({ grep '^{' "$TERM_LOG" || true; } \ - | jq -r 'select(.event == "exiting") | .reason' 2>/dev/null | tail -1 || true) + # Each line on its own: a line that is not exactly one JSON object (an echoed + # token, or the event with a token glued on) must not count. + TERM_REASON=$(jq -rR 'fromjson? | select(type == "object" and .event == "exiting") | .reason' \ + "$TERM_LOG" 2>/dev/null | tail -1 || true) if [ "$TERM_HUNG" = true ]; then fail "SIGTERM mid-stream: server still running 30s after SIGTERM" - tail -5 "$TERM_LOG" || true + tail -40 "$TERM_LOG" || true + cp "$TERM_LOG" /tmp/SwiftLM-test-sigterm.log 2>/dev/null || true elif [ "$TERM_STATUS" -eq 0 ] && [ "$TERM_REASON" = "requested" ]; then pass "SIGTERM mid-stream: exit 0, exiting{requested} on its own line" else fail "SIGTERM mid-stream: status=$TERM_STATUS reason='${TERM_REASON}'" - tail -5 "$TERM_LOG" || true + tail -40 "$TERM_LOG" || true + cp "$TERM_LOG" /tmp/SwiftLM-test-sigterm.log 2>/dev/null || true fi fi rm -f "$TERM_LOG" "$TERM_STREAM" +unset TERM_LOG TERM_STREAM # ── Results ────────────────────────────────────────────────────────── echo "" From f96e756c8f98c26b82fc9f5369bac85548d5c594 Mon Sep 17 00:00:00 2001 From: simba Date: Sun, 27 Sep 2026 21:56:00 -0700 Subject: [PATCH 5/7] fix(ssd): use loadable local copies first; consult the Hub only to download Review round 3 for #198. Checking every start against the Hub's file list rejected copies the loader reads fine, contacted the Hub on starts main made offline, and failed on any Hub error. Back to local-first: - Non-streaming starts are unchanged from main: no Hub request. - --stream-experts uses the first local copy the loader can read (config.json, tokenizer.json, and either every indexed shard or any top-level *.safetensors) without contacting the Hub. Only when none exists does it ask the Hub and download into the Application Support copy. - Hub 401/404 (with no loadable copy) is ModelNotOnHub; any other Hub or network error is ModelUnavailableOffline. No custom timeout. - The listing only checks the repo has tokenizer.json and which top-level files the finished download must contain. - An in-progress marker is written before a download into Application Support and removed after success, so an interrupted download or update is re-fetched instead of trusted. Co-Authored-By: Claude Opus 5.5 --- Sources/MLXInferenceCore/ModelStorage.swift | 5 - Sources/SwiftLM/Server.swift | 26 +-- Sources/SwiftLM/StreamingDirectory.swift | 191 +++++++++----------- 3 files changed, 93 insertions(+), 129 deletions(-) diff --git a/Sources/MLXInferenceCore/ModelStorage.swift b/Sources/MLXInferenceCore/ModelStorage.swift index da7429e..9c39c85 100644 --- a/Sources/MLXInferenceCore/ModelStorage.swift +++ b/Sources/MLXInferenceCore/ModelStorage.swift @@ -280,11 +280,6 @@ public enum ModelStorage { validateModelFiles(in: directory, logFailures: true) } - /// `validateLocalModelDirectory` without logging, for probing candidate copies. - public static func validateLocalModelDirectory(_ directory: URL, logFailures: Bool) -> Bool { - validateModelFiles(in: directory, logFailures: logFailures) - } - /// Read the raw config.json dictionary for a downloaded model. /// Verifies that all required safetensors files are present in the snapshot directory. /// This prevents the engine from entering `.ready` state if a download was interrupted or corrupted. diff --git a/Sources/SwiftLM/Server.swift b/Sources/SwiftLM/Server.swift index 3c097d8..1f252d9 100644 --- a/Sources/SwiftLM/Server.swift +++ b/Sources/SwiftLM/Server.swift @@ -674,30 +674,12 @@ struct MLXServer: AsyncParsableCommand { .appendingPathComponent("HuggingFace", isDirectory: true)) let isHubId = !ModelStorage.isLocalDirectoryPath(modelId) && !FileManager.default.fileExists(atPath: modelId) - let validatedLocal = isHubId ? ModelStorage.validatedContentDirectory(for: modelId) : nil - // A copy whose download stopped after the shards still validates, so check local - // copies against the Hub's file list. --stream-experts needs the list anyway, and - // a Hub error there (unknown or gated id) fails the start. - var hubListing: HubListing? = nil - if isHubId, !self.info, self.streamExperts || validatedLocal != nil { - phase = .architectureProbe - hubListing = self.streamExperts - ? try await fetchHubListing(cliHub, modelId: modelId) - : try? await fetchHubListing(cliHub, modelId: modelId) - } - var validatedLocalMissing: [String] = [] - if let dir = validatedLocal, case .files(let files)? = hubListing { - validatedLocalMissing = missingFiles(files, in: dir) - } var modelConfig: ModelConfiguration if ModelStorage.isLocalDirectoryPath(modelId) { print("[SwiftLM] Loading from local directory: \(modelId)") modelConfig = ModelConfiguration(directory: URL(filePath: modelId)) - } else if let dir = validatedLocal, !validatedLocalMissing.isEmpty { - print("[SwiftLM] \(dir.path) is missing \(validatedLocalMissing.count) file(s) (e.g. \(validatedLocalMissing[0])); loading through the Hub to complete it.") - modelConfig = ModelConfiguration(id: modelId) - } else if let localDirectory = validatedLocal { + } else if let localDirectory = ModelStorage.validatedContentDirectory(for: modelId) { // Any validated copy in the shared HF cache, in any supported layout. Note // this deliberately does NOT use localLoadDirectory: that skips the // materialized `models//` layout on the grounds that HubApi @@ -723,9 +705,11 @@ struct MLXServer: AsyncParsableCommand { var modelDirectory = ModelStorage.validatedContentDirectory(for: modelId) ?? resolveModelDirectory(modelId: modelId) - if self.streamExperts, !self.info, isHubId, let listing = hubListing { + if self.streamExperts, !self.info, isHubId { + // A Hub or download failure here is a model problem, not a binary one. + phase = .architectureProbe modelDirectory = try await resolveStreamingDirectory( - modelId: modelId, candidate: modelDirectory, hub: cliHub, listing: listing) + modelId: modelId, candidate: modelDirectory, hub: cliHub) } var mainModelProfile: ModelProfile? = nil if self.streamExperts, let dir = modelDirectory { diff --git a/Sources/SwiftLM/StreamingDirectory.swift b/Sources/SwiftLM/StreamingDirectory.swift index 63ee5a3..37eb2ca 100644 --- a/Sources/SwiftLM/StreamingDirectory.swift +++ b/Sources/SwiftLM/StreamingDirectory.swift @@ -1,10 +1,8 @@ -// StreamingDirectory.swift — where a Hub model is loaded from, checked against the -// Hub's own file list. +// StreamingDirectory.swift — the directory `--stream-experts` loads a Hub model from. // -// `--stream-experts` streams from the one directory it was activated for, and a -// directory load doesn't fill in missing files. So the directory must be complete, -// and "complete" is decided by what the Hub lists, not by guessing a shard layout -// (repos ship `weights.NN.safetensors`, no index, or an index naming extra shards). +// Streaming targets one directory, and a directory load doesn't fill in missing +// files, so that copy must be loadable. Local copies are used as they are when they +// have what the loader reads; the Hub is consulted only when none does. import Foundation import Hub @@ -13,7 +11,7 @@ import MLXInferenceCore /// Files SwiftLM downloads for a model: weights, configs, tokenizer and templates. let modelDownloadPatterns = ["*.safetensors", "*.json", "*.jinja"] -/// A `--stream-experts` prefetch finished without every listed file on disk. +/// A `--stream-experts` download finished without the files the loader needs. struct ModelDownloadIncomplete: Error, CustomStringConvertible { let modelId: String let directory: URL @@ -30,128 +28,115 @@ struct ModelMissingTokenizer: Error, CustomStringConvertible { } } -/// The Hub couldn't be reached and there is no local copy known to be complete. -struct ModelUnavailableOffline: Error, CustomStringConvertible { +/// The Hub answered 401/404 and there is no loadable local copy. +struct ModelNotOnHub: Error, CustomStringConvertible { let modelId: String let underlying: Error var description: String { - "Could not reach the Hub to fetch \(modelId) (\(underlying)), and no complete local copy was found." + "The Hub has no accessible repository \(modelId) (\(underlying)), and no loadable local copy was found. Check the id, or set HF_TOKEN for a private repository." } } -/// The Hub answered, but has no repository `modelId` this user can read. -struct ModelNotOnHub: Error, CustomStringConvertible { +/// The Hub couldn't be used and there is no loadable local copy. +struct ModelUnavailableOffline: Error, CustomStringConvertible { let modelId: String let underlying: Error var description: String { - "The Hub has no accessible repository \(modelId) (\(underlying)). Check the id, or log in for a gated model." + "Could not get \(modelId) from the Hub (\(underlying)), and no loadable local copy was found." } } -struct HubListingTimeout: Error, CustomStringConvertible { - var description: String { "the Hub did not answer in time" } -} - -/// The Hub's file list for a repository, or why there isn't one. -enum HubListing { - case files([String]) - /// A network error or timeout; the Hub gave no answer. - case unreachable(Error) +/// Whether `directory` has `tokenizer.json`. swift-transformers requires it; +/// `tokenizer_config.json` is optional, and `vocab.json` is no substitute. +func hasTokenizerJSON(in directory: URL) -> Bool { + FileManager.default.fileExists(atPath: directory.appendingPathComponent("tokenizer.json").path) } -/// Asks the Hub which files `modelId` has. Throws only when the Hub answered with -/// an error (unknown or gated repository); connectivity problems are `.unreachable`. -func fetchHubListing(_ hub: HubApi, modelId: String) async throws -> HubListing { - do { - let files = try await withThrowingTaskGroup(of: [String].self) { group in - group.addTask { try await hub.getFilenames(from: modelId, matching: modelDownloadPatterns) } - group.addTask { - try await Task.sleep(nanoseconds: 15_000_000_000) - throw HubListingTimeout() - } - defer { group.cancelAll() } - return try await group.next()! - } - return .files(files) - } catch let error as Hub.HubClientError { - switch error { - case .httpStatusCode, .authorizationRequired, .resourceNotFound, .fileNotFound: - throw ModelNotOnHub(modelId: modelId, underlying: error) - default: - return .unreachable(error) - } - } catch { - return .unreachable(error) +/// Whether the loader can read `directory` as a model: `config.json`, `tokenizer.json`, +/// and weights. With an index, every shard it names; without one, any top-level +/// `*.safetensors` (`model.safetensors`, `weights.00.safetensors`, ...). +func hasLoadableModelFiles(in directory: URL) -> Bool { + let fm = FileManager.default + func nonEmpty(_ name: String) -> Bool { + let url = directory.appendingPathComponent(name).resolvingSymlinksInPath() + let size = (try? fm.attributesOfItem(atPath: url.path)[.size] as? Int) ?? 0 + return size > 0 } + guard nonEmpty("config.json"), nonEmpty("tokenizer.json") else { return false } + let index = directory.appendingPathComponent("model.safetensors.index.json") + if let data = try? Data(contentsOf: index), + let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any], + let weightMap = json["weight_map"] as? [String: String] + { + return Set(weightMap.values).allSatisfy(nonEmpty) + } + let names = (try? fm.contentsOfDirectory(atPath: directory.path)) ?? [] + return names.contains { $0.hasSuffix(".safetensors") && nonEmpty($0) } } -/// Listed files that are absent from `directory`. -func missingFiles(_ files: [String], in directory: URL) -> [String] { - files.filter { !FileManager.default.fileExists(atPath: directory.appendingPathComponent($0).path) } -} - -/// Whether `directory` has `tokenizer.json`. swift-transformers requires it; -/// `tokenizer_config.json` is optional, and `vocab.json` is no substitute. -func hasTokenizerJSON(in directory: URL) -> Bool { - FileManager.default.fileExists(atPath: directory.appendingPathComponent("tokenizer.json").path) +/// Whether a Hub error means the repository isn't available to this user. +private func isNotOnHub(_ error: Error) -> Bool { + guard let hubError = error as? Hub.HubClientError else { return false } + switch hubError { + case .authorizationRequired, .resourceNotFound: return true + case .httpStatusCode(let code): return code == 401 || code == 404 + default: return false + } } /// The directory `--stream-experts` loads and streams `modelId` from. /// -/// Online, a local copy (`candidate`, then the loader's Application Support copy) -/// is used only if it has every listed file; otherwise the model is downloaded and -/// must then have them all. Offline, only a copy known to be complete is used. -func resolveStreamingDirectory( - modelId: String, candidate: URL?, hub: HubApi, listing: HubListing -) async throws -> URL { +/// A loadable local copy (`candidate`, then the loader's Application Support copy) +/// is used without contacting the Hub. Otherwise the model is downloaded into the +/// Application Support copy, which must then have every top-level listed file. +func resolveStreamingDirectory(modelId: String, candidate: URL?, hub: HubApi) async throws -> URL { let fm = FileManager.default let localRepo = hub.localRepoLocation(Hub.Repo(id: modelId)) - // Records that this copy once matched the Hub listing, for offline starts. - let marker = localRepo.appendingPathComponent(".swiftlm-snapshot-complete") + // Present while a download into localRepo is unfinished; that copy isn't trusted. + let inProgress = localRepo.appendingPathComponent(".swiftlm-download-in-progress") let copies = [candidate, localRepo].compactMap { $0 }.filter { fm.fileExists(atPath: $0.path) } + let loadable = copies.filter { + hasLoadableModelFiles(in: $0) && !($0 == localRepo && fm.fileExists(atPath: inProgress.path)) + } + if let dir = loadable.first { return dir } - switch listing { - case .files(let files): - guard files.contains("tokenizer.json") else { throw ModelMissingTokenizer(modelId: modelId) } - for dir in copies { - let missing = missingFiles(files, in: dir) - if missing.isEmpty { - if dir == localRepo { fm.createFile(atPath: marker.path, contents: nil) } - return dir - } - print("[SwiftLM] \(dir.path) is missing \(missing.count) of \(files.count) files (e.g. \(missing[0])).") - } - // A dense model won't stream: leave completing it to the loader instead of - // downloading everything here first. - if let dir = copies.first(where: { - fm.fileExists(atPath: $0.appendingPathComponent("config.json").path) - }), let profile = ModelProfiler.profile(modelDirectory: dir, modelId: modelId), !profile.isMoE { - return dir - } - - print("[SwiftLM] --stream-experts: downloading \(modelId) before loading...") - let tracker = ProgressTracker(modelId: modelId) - defer { tracker.finish() } - let snapshot = try await hub.snapshot(from: modelId, matching: modelDownloadPatterns) { - tracker.printProgress($0) - } - tracker.finish() - guard missingFiles(files, in: snapshot).isEmpty else { - throw ModelDownloadIncomplete(modelId: modelId, directory: snapshot) - } - fm.createFile(atPath: marker.path, contents: nil) - return snapshot + let files: [String] + do { + files = try await hub.getFilenames(from: modelId, matching: modelDownloadPatterns) + } catch where isNotOnHub(error) { + throw ModelNotOnHub(modelId: modelId, underlying: error) + } catch { + throw ModelUnavailableOffline(modelId: modelId, underlying: error) + } + guard files.contains("tokenizer.json") else { throw ModelMissingTokenizer(modelId: modelId) } - case .unreachable(let error): - let verified = fm.fileExists(atPath: marker.path) && hasTokenizerJSON(in: localRepo) - ? [localRepo] : [] - let plausible = copies.filter { - hasTokenizerJSON(in: $0) && ModelStorage.validateLocalModelDirectory($0, logFailures: false) - } - guard let dir = (verified + plausible).first else { - throw ModelUnavailableOffline(modelId: modelId, underlying: error) - } - print("[SwiftLM] ⚠️ Could not reach the Hub to check \(modelId) (\(error)); using \(dir.path).") + // A dense model won't stream: leave completing it to the loader instead of + // downloading everything here first. + if let dir = copies.first(where: { + fm.fileExists(atPath: $0.appendingPathComponent("config.json").path) + }), let profile = ModelProfiler.profile(modelDirectory: dir, modelId: modelId), !profile.isMoE { return dir } + for dir in copies { + let reason = dir == localRepo && fm.fileExists(atPath: inProgress.path) + ? "its download didn't finish" : "missing weights or tokenizer.json" + print("[SwiftLM] \(dir.path) can't be loaded as is (\(reason)).") + } + + print("[SwiftLM] --stream-experts: downloading \(modelId) before loading...") + try fm.createDirectory(at: localRepo, withIntermediateDirectories: true) + fm.createFile(atPath: inProgress.path, contents: nil) + let tracker = ProgressTracker(modelId: modelId) + defer { tracker.finish() } + let snapshot = try await hub.snapshot(from: modelId, matching: modelDownloadPatterns) { + tracker.printProgress($0) + } + tracker.finish() + let topLevel = files.filter { !$0.contains("/") } + let missing = topLevel.filter { !fm.fileExists(atPath: snapshot.appendingPathComponent($0).path) } + guard missing.isEmpty, hasLoadableModelFiles(in: snapshot) else { + throw ModelDownloadIncomplete(modelId: modelId, directory: snapshot) + } + try? fm.removeItem(at: inProgress) + return snapshot } From cdb16e6e1666bafb819b04c3fd862428e0fbac11 Mon Sep 17 00:00:00 2001 From: simba Date: Sun, 27 Sep 2026 22:32:03 -0700 Subject: [PATCH 6/7] fix(ssd): judge local copies the way the loader reads them Review round 4 for #198: - localWeightState mirrors safetensorWeightURLs: the index only when every file it names exists, otherwise model*, weight*, then all top-level *.safetensors. Shards named stem-NNNNN-of-MMMMM must all be present, a lone model.safetensors is complete, and other layouts are unverified, so the Hub listing decides. A partial download without an index is no longer served, and a complete copy with a stale index (Qwen3-VL MoE repos) is no longer re-downloaded forever. - Unverified copies are checked against the listing; if the Hub can't be used they are used with a warning. - A Hub 404 surfaces as .fileNotFound in swift-transformers; treat it as not on the Hub. - A dangling symlink counts as missing. - resolveStreamingDirectory reports whether the copy is complete, so a dense model's complete copy loads by directory after streaming turns off (no network), and only the dense shortcut's partial copy loads by id. - Unit tests: partial no-index, complete and partial stale-index, single file, unverified layout, dangling symlink, tokenizer.json required. Co-Authored-By: Claude Opus 5.5 --- Sources/SwiftLM/Server.swift | 8 +- Sources/SwiftLM/StreamingDirectory.swift | 163 +++++++++++++----- .../StreamingDirectoryTests.swift | 88 ++++++++++ 3 files changed, 210 insertions(+), 49 deletions(-) create mode 100644 tests/SwiftLMTests/StreamingDirectoryTests.swift diff --git a/Sources/SwiftLM/Server.swift b/Sources/SwiftLM/Server.swift index 1f252d9..157c425 100644 --- a/Sources/SwiftLM/Server.swift +++ b/Sources/SwiftLM/Server.swift @@ -705,11 +705,14 @@ struct MLXServer: AsyncParsableCommand { var modelDirectory = ModelStorage.validatedContentDirectory(for: modelId) ?? resolveModelDirectory(modelId: modelId) + var modelDirectoryComplete = false if self.streamExperts, !self.info, isHubId { // A Hub or download failure here is a model problem, not a binary one. phase = .architectureProbe - modelDirectory = try await resolveStreamingDirectory( + let resolved = try await resolveStreamingDirectory( modelId: modelId, candidate: modelDirectory, hub: cliHub) + modelDirectory = resolved.directory + modelDirectoryComplete = resolved.complete } var mainModelProfile: ModelProfile? = nil if self.streamExperts, let dir = modelDirectory { @@ -730,7 +733,8 @@ struct MLXServer: AsyncParsableCommand { // Streaming is activated for `modelDirectory`, and only a load of that exact // directory streams. Load from it, or a different lookup could pick another copy. - if self.streamExperts, let dir = modelDirectory { + // A complete copy also loads from here after the MoE check turns streaming off. + if let dir = modelDirectory, self.streamExperts || modelDirectoryComplete { modelConfig = ModelConfiguration(directory: dir) } diff --git a/Sources/SwiftLM/StreamingDirectory.swift b/Sources/SwiftLM/StreamingDirectory.swift index 37eb2ca..022260b 100644 --- a/Sources/SwiftLM/StreamingDirectory.swift +++ b/Sources/SwiftLM/StreamingDirectory.swift @@ -1,8 +1,8 @@ // StreamingDirectory.swift — the directory `--stream-experts` loads a Hub model from. // // Streaming targets one directory, and a directory load doesn't fill in missing -// files, so that copy must be loadable. Local copies are used as they are when they -// have what the loader reads; the Hub is consulted only when none does. +// files, so that copy must be complete. A local copy whose files show it is complete +// is used as is; otherwise the Hub listing decides. import Foundation import Hub @@ -28,99 +28,168 @@ struct ModelMissingTokenizer: Error, CustomStringConvertible { } } -/// The Hub answered 401/404 and there is no loadable local copy. +/// The Hub answered 401/404 and there is no usable local copy. struct ModelNotOnHub: Error, CustomStringConvertible { let modelId: String let underlying: Error var description: String { - "The Hub has no accessible repository \(modelId) (\(underlying)), and no loadable local copy was found. Check the id, or set HF_TOKEN for a private repository." + "The Hub has no accessible repository \(modelId) (\(underlying)), and no usable local copy was found. Check the id, or set HF_TOKEN for a private repository." } } -/// The Hub couldn't be used and there is no loadable local copy. +/// The Hub couldn't be used and there is no usable local copy. struct ModelUnavailableOffline: Error, CustomStringConvertible { let modelId: String let underlying: Error var description: String { - "Could not get \(modelId) from the Hub (\(underlying)), and no loadable local copy was found." + "Could not get \(modelId) from the Hub (\(underlying)), and no usable local copy was found." } } -/// Whether `directory` has `tokenizer.json`. swift-transformers requires it; -/// `tokenizer_config.json` is optional, and `vocab.json` is no substitute. -func hasTokenizerJSON(in directory: URL) -> Bool { - FileManager.default.fileExists(atPath: directory.appendingPathComponent("tokenizer.json").path) +/// Whether a local copy's weights are known to be complete. +enum LocalWeights: Equatable { + case complete + case incomplete + /// A layout the file names can't vouch for (`weights.00.safetensors`, say); + /// the Hub listing decides. + case unverified } -/// Whether the loader can read `directory` as a model: `config.json`, `tokenizer.json`, -/// and weights. With an index, every shard it names; without one, any top-level -/// `*.safetensors` (`model.safetensors`, `weights.00.safetensors`, ...). -func hasLoadableModelFiles(in directory: URL) -> Bool { +/// Whether `name` in `directory` exists and is non-empty, following symlinks +/// (Hub snapshots link to blobs, and a dangling link counts as missing). +func isPresentFile(_ name: String, in directory: URL) -> Bool { + let url = directory.appendingPathComponent(name) let fm = FileManager.default - func nonEmpty(_ name: String) -> Bool { - let url = directory.appendingPathComponent(name).resolvingSymlinksInPath() - let size = (try? fm.attributesOfItem(atPath: url.path)[.size] as? Int) ?? 0 - return size > 0 - } - guard nonEmpty("config.json"), nonEmpty("tokenizer.json") else { return false } + guard fm.fileExists(atPath: url.path) else { return false } + let attributes = try? fm.attributesOfItem(atPath: url.resolvingSymlinksInPath().path) + return ((attributes?[.size] as? Int) ?? 0) > 0 +} + +/// `stem-00002-of-00004.safetensors` as (stem, 2, 4); nil for other names. +func shardNumber(_ name: String) -> (stem: String, index: Int, count: Int)? { + guard name.hasSuffix(".safetensors") else { return nil } + let base = name.dropLast(".safetensors".count) + guard let of = base.range(of: "-of-", options: .backwards) else { return nil } + let countText = base[of.upperBound...] + let head = base[.. LocalWeights { let index = directory.appendingPathComponent("model.safetensors.index.json") if let data = try? Data(contentsOf: index), let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any], - let weightMap = json["weight_map"] as? [String: String] + let weightMap = json["weight_map"] as? [String: String], + !weightMap.isEmpty, + Set(weightMap.values).allSatisfy({ isPresentFile($0, in: directory) }) { - return Set(weightMap.values).allSatisfy(nonEmpty) + return .complete + } + let top = ((try? FileManager.default.contentsOfDirectory(atPath: directory.path)) ?? []) + .filter { $0.hasSuffix(".safetensors") } + let chosen = ["model", "weight"] + .map { prefix in top.filter { $0.hasPrefix(prefix) } } + .first { !$0.isEmpty } ?? top + guard !chosen.isEmpty else { return .incomplete } + + // `stem-00002-of-00004.safetensors`: every shard 1...4 must be there. + var groups = [String: (count: Int, have: Set)]() + for name in chosen { + guard let (stem, i, n) = shardNumber(name) else { continue } + let key = "\(stem)/\(n)" + var group = groups[key] ?? (n, []) + if isPresentFile(name, in: directory) { group.have.insert(i) } + groups[key] = group + } + if !groups.isEmpty { + let whole = groups.values.allSatisfy { $0.count > 0 && $0.have == Set(1 ... $0.count) } + return whole ? .complete : .incomplete } - let names = (try? fm.contentsOfDirectory(atPath: directory.path)) ?? [] - return names.contains { $0.hasSuffix(".safetensors") && nonEmpty($0) } + if chosen == ["model.safetensors"] { + return isPresentFile("model.safetensors", in: directory) ? .complete : .incomplete + } + return .unverified +} + +/// Whether `directory` has the non-weight files the loader reads. swift-transformers +/// requires `tokenizer.json`; `tokenizer_config.json` is optional. +func hasModelConfigAndTokenizer(in directory: URL) -> Bool { + isPresentFile("config.json", in: directory) && isPresentFile("tokenizer.json", in: directory) } /// Whether a Hub error means the repository isn't available to this user. private func isNotOnHub(_ error: Error) -> Bool { guard let hubError = error as? Hub.HubClientError else { return false } switch hubError { - case .authorizationRequired, .resourceNotFound: return true + case .authorizationRequired, .resourceNotFound, .fileNotFound: return true case .httpStatusCode(let code): return code == 401 || code == 404 default: return false } } -/// The directory `--stream-experts` loads and streams `modelId` from. +/// The directory `--stream-experts` loads and streams `modelId` from, and whether that +/// copy is complete (a dense model's partial copy is returned for the loader to finish). /// -/// A loadable local copy (`candidate`, then the loader's Application Support copy) -/// is used without contacting the Hub. Otherwise the model is downloaded into the -/// Application Support copy, which must then have every top-level listed file. -func resolveStreamingDirectory(modelId: String, candidate: URL?, hub: HubApi) async throws -> URL { +/// A local copy (`candidate`, then the loader's Application Support copy) whose files +/// show it is complete is used without contacting the Hub. Otherwise the Hub listing +/// decides: a copy with every top-level listed file is used, or the model is downloaded +/// into the Application Support copy. If the Hub can't be used, a copy whose layout +/// its names can't verify is used with a warning. +func resolveStreamingDirectory( + modelId: String, candidate: URL?, hub: HubApi +) async throws -> (directory: URL, complete: Bool) { let fm = FileManager.default let localRepo = hub.localRepoLocation(Hub.Repo(id: modelId)) // Present while a download into localRepo is unfinished; that copy isn't trusted. let inProgress = localRepo.appendingPathComponent(".swiftlm-download-in-progress") - let copies = [candidate, localRepo].compactMap { $0 }.filter { fm.fileExists(atPath: $0.path) } - let loadable = copies.filter { - hasLoadableModelFiles(in: $0) && !($0 == localRepo && fm.fileExists(atPath: inProgress.path)) + let copies = [candidate, localRepo].compactMap { $0 }.filter { + fm.fileExists(atPath: $0.path) + && !($0 == localRepo && fm.fileExists(atPath: inProgress.path)) + } + let usable = copies.filter(hasModelConfigAndTokenizer) + if let dir = usable.first(where: { localWeightState(in: $0) == .complete }) { + return (dir, true) } - if let dir = loadable.first { return dir } + let unverified = usable.filter { localWeightState(in: $0) == .unverified } let files: [String] do { files = try await hub.getFilenames(from: modelId, matching: modelDownloadPatterns) - } catch where isNotOnHub(error) { - throw ModelNotOnHub(modelId: modelId, underlying: error) } catch { + if let dir = unverified.first { + print("[SwiftLM] ⚠️ Could not check \(modelId) with the Hub (\(error)); using \(dir.path) as is.") + return (dir, true) + } + if isNotOnHub(error) { throw ModelNotOnHub(modelId: modelId, underlying: error) } throw ModelUnavailableOffline(modelId: modelId, underlying: error) } guard files.contains("tokenizer.json") else { throw ModelMissingTokenizer(modelId: modelId) } + let topLevel = files.filter { !$0.contains("/") } + if let dir = usable.first(where: { dir in topLevel.allSatisfy { isPresentFile($0, in: dir) } }) { + return (dir, true) + } // A dense model won't stream: leave completing it to the loader instead of // downloading everything here first. - if let dir = copies.first(where: { - fm.fileExists(atPath: $0.appendingPathComponent("config.json").path) - }), let profile = ModelProfiler.profile(modelDirectory: dir, modelId: modelId), !profile.isMoE { - return dir + let allCopies = [candidate, localRepo].compactMap { $0 } + if let dir = allCopies.first(where: { isPresentFile("config.json", in: $0) }), + let profile = ModelProfiler.profile(modelDirectory: dir, modelId: modelId), !profile.isMoE + { + return (dir, false) } - for dir in copies { + for dir in allCopies where fm.fileExists(atPath: dir.path) { let reason = dir == localRepo && fm.fileExists(atPath: inProgress.path) - ? "its download didn't finish" : "missing weights or tokenizer.json" - print("[SwiftLM] \(dir.path) can't be loaded as is (\(reason)).") + ? "its download didn't finish" : "it is missing files" + print("[SwiftLM] \(dir.path) can't be used as is (\(reason)).") } print("[SwiftLM] --stream-experts: downloading \(modelId) before loading...") @@ -132,11 +201,11 @@ func resolveStreamingDirectory(modelId: String, candidate: URL?, hub: HubApi) as tracker.printProgress($0) } tracker.finish() - let topLevel = files.filter { !$0.contains("/") } - let missing = topLevel.filter { !fm.fileExists(atPath: snapshot.appendingPathComponent($0).path) } - guard missing.isEmpty, hasLoadableModelFiles(in: snapshot) else { + guard topLevel.allSatisfy({ isPresentFile($0, in: snapshot) }), + hasModelConfigAndTokenizer(in: snapshot) + else { throw ModelDownloadIncomplete(modelId: modelId, directory: snapshot) } try? fm.removeItem(at: inProgress) - return snapshot + return (snapshot, true) } diff --git a/tests/SwiftLMTests/StreamingDirectoryTests.swift b/tests/SwiftLMTests/StreamingDirectoryTests.swift new file mode 100644 index 0000000..d555cfb --- /dev/null +++ b/tests/SwiftLMTests/StreamingDirectoryTests.swift @@ -0,0 +1,88 @@ +import Foundation +import XCTest + +@testable import SwiftLM + +/// `localWeightState` must judge a copy the way the loader reads it, so a partial +/// download is never served and a complete copy with a stale index is not rejected. +final class StreamingDirectoryTests: XCTestCase { + private var dir: URL! + + override func setUpWithError() throws { + dir = FileManager.default.temporaryDirectory + .appendingPathComponent("streaming-dir-\(UUID().uuidString)") + try FileManager.default.createDirectory(at: dir, withIntermediateDirectories: true) + try write("config.json") + try write("tokenizer.json") + } + + override func tearDownWithError() throws { + try? FileManager.default.removeItem(at: dir) + } + + private func write(_ name: String, _ text: String = "{}") throws { + try Data(text.utf8).write(to: dir.appendingPathComponent(name)) + } + + private func shards(_ indices: [Int], of count: Int) throws { + for i in indices { + try write(String(format: "model-%05d-of-%05d.safetensors", i, count), "weights") + } + } + + /// An index naming shards the repository doesn't ship (carried over from the source). + private func staleIndex(naming count: Int) throws { + let map = Dictionary(uniqueKeysWithValues: (1 ... count).map { + ("layer\($0).weight", String(format: "model-%05d-of-%05d.safetensors", $0, count)) + }) + let json = try JSONSerialization.data(withJSONObject: ["weight_map": map]) + try json.write(to: dir.appendingPathComponent("model.safetensors.index.json")) + } + + func testPartialCopyWithoutIndexIsIncomplete() throws { + try shards([1, 3], of: 3) + XCTAssertEqual(localWeightState(in: dir), .incomplete) + } + + func testAllShardsWithoutIndexAreComplete() throws { + try shards([1, 2, 3], of: 3) + XCTAssertEqual(localWeightState(in: dir), .complete) + } + + func testStaleIndexWithEveryShippedShardIsComplete() throws { + try staleIndex(naming: 13) + try shards([1, 2, 3, 4], of: 4) + XCTAssertEqual(localWeightState(in: dir), .complete) + } + + func testStaleIndexWithAMissingShardIsIncomplete() throws { + try staleIndex(naming: 13) + try shards([1, 2, 4], of: 4) + XCTAssertEqual(localWeightState(in: dir), .incomplete) + } + + func testSingleModelFileIsComplete() throws { + try write("model.safetensors", "weights") + XCTAssertEqual(localWeightState(in: dir), .complete) + } + + func testUnconventionalLayoutIsUnverified() throws { + try write("weights.00.safetensors", "weights") + XCTAssertEqual(localWeightState(in: dir), .unverified) + } + + func testDanglingSymlinkIsMissing() throws { + try FileManager.default.createSymbolicLink( + at: dir.appendingPathComponent("model.safetensors"), + withDestinationURL: dir.appendingPathComponent("gone.bin")) + XCTAssertFalse(isPresentFile("model.safetensors", in: dir)) + XCTAssertEqual(localWeightState(in: dir), .incomplete) + } + + func testTokenizerJSONIsRequired() throws { + XCTAssertTrue(hasModelConfigAndTokenizer(in: dir)) + try FileManager.default.removeItem(at: dir.appendingPathComponent("tokenizer.json")) + try write("tokenizer_config.json") + XCTAssertFalse(hasModelConfigAndTokenizer(in: dir)) + } +} From ca32d58c4445791fb04a1a4b038cc9ba609f0d62 Mon Sep 17 00:00:00 2001 From: simba Date: Sun, 27 Sep 2026 23:06:29 -0700 Subject: [PATCH 7/7] fix(ssd): a missing indexed file in a subdirectory means incomplete MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review round 5 for #198: - localWeightState falls back to the shard check only when the index's missing names are top-level (a stale index). A missing file in a subdirectory, such as OptiQ's optiq/optiq_vision.safetensors, makes the copy incomplete; before, a model*-only download of an OptiQ VLM was served and its vision tower silently ran on random weights. - The Hub-listing and post-download checks also require listed nested files the copy's index names (requiredListedFiles). - --info --stream-experts drops a clearly incomplete candidate again, so it says the model isn't downloaded instead of profiling a partial copy. - Unit tests: nested indexed file missing → incomplete; present → complete; requiredListedFiles keeps indexed nested files and skips unindexed ones. Co-Authored-By: Claude Opus 5.5 --- Sources/SwiftLM/Server.swift | 6 +++ Sources/SwiftLM/StreamingDirectory.swift | 39 +++++++++++++------ .../StreamingDirectoryTests.swift | 35 +++++++++++++++++ 3 files changed, 69 insertions(+), 11 deletions(-) diff --git a/Sources/SwiftLM/Server.swift b/Sources/SwiftLM/Server.swift index 157c425..bd7225d 100644 --- a/Sources/SwiftLM/Server.swift +++ b/Sources/SwiftLM/Server.swift @@ -706,6 +706,12 @@ struct MLXServer: AsyncParsableCommand { ModelStorage.validatedContentDirectory(for: modelId) ?? resolveModelDirectory(modelId: modelId) var modelDirectoryComplete = false + // --info doesn't download; don't profile a copy that is clearly partial. + if self.streamExperts, self.info, isHubId, let dir = modelDirectory, + localWeightState(in: dir) == .incomplete + { + modelDirectory = nil + } if self.streamExperts, !self.info, isHubId { // A Hub or download failure here is a model problem, not a binary one. phase = .architectureProbe diff --git a/Sources/SwiftLM/StreamingDirectory.swift b/Sources/SwiftLM/StreamingDirectory.swift index 022260b..dab5c3e 100644 --- a/Sources/SwiftLM/StreamingDirectory.swift +++ b/Sources/SwiftLM/StreamingDirectory.swift @@ -85,14 +85,13 @@ func shardNumber(_ name: String) -> (stem: String, index: Int, count: Int)? { /// `safetensorWeightURLs`: the index when every file it names exists, otherwise /// `model*`, then `weight*`, then every top-level `*.safetensors`. func localWeightState(in directory: URL) -> LocalWeights { - let index = directory.appendingPathComponent("model.safetensors.index.json") - if let data = try? Data(contentsOf: index), - let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any], - let weightMap = json["weight_map"] as? [String: String], - !weightMap.isEmpty, - Set(weightMap.values).allSatisfy({ isPresentFile($0, in: directory) }) - { - return .complete + let indexed = indexedWeightFiles(in: directory) + if !indexed.isEmpty { + let missing = indexed.filter { !isPresentFile($0, in: directory) } + if missing.isEmpty { return .complete } + // A stale index names top-level shards the repo doesn't ship; a missing file + // in a subdirectory (optiq/optiq_vision.safetensors) is really missing. + if missing.contains(where: { $0.contains("/") }) { return .incomplete } } let top = ((try? FileManager.default.contentsOfDirectory(atPath: directory.path)) ?? []) .filter { $0.hasSuffix(".safetensors") } @@ -120,6 +119,23 @@ func localWeightState(in directory: URL) -> LocalWeights { return .unverified } +/// The files `model.safetensors.index.json` in `directory` names, if any. +func indexedWeightFiles(in directory: URL) -> Set { + let index = directory.appendingPathComponent("model.safetensors.index.json") + guard let data = try? Data(contentsOf: index), + let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any], + let weightMap = json["weight_map"] as? [String: String] + else { return [] } + return Set(weightMap.values) +} + +/// Listed files `directory` must have: every top-level one, plus any nested one its +/// index names (weights in a subdirectory, such as OptiQ's vision tower). +func requiredListedFiles(_ files: [String], in directory: URL) -> [String] { + let indexed = indexedWeightFiles(in: directory) + return files.filter { !$0.contains("/") || indexed.contains($0) } +} + /// Whether `directory` has the non-weight files the loader reads. swift-transformers /// requires `tokenizer.json`; `tokenizer_config.json` is optional. func hasModelConfigAndTokenizer(in directory: URL) -> Bool { @@ -173,8 +189,9 @@ func resolveStreamingDirectory( throw ModelUnavailableOffline(modelId: modelId, underlying: error) } guard files.contains("tokenizer.json") else { throw ModelMissingTokenizer(modelId: modelId) } - let topLevel = files.filter { !$0.contains("/") } - if let dir = usable.first(where: { dir in topLevel.allSatisfy { isPresentFile($0, in: dir) } }) { + if let dir = usable.first(where: { dir in + requiredListedFiles(files, in: dir).allSatisfy { isPresentFile($0, in: dir) } + }) { return (dir, true) } @@ -201,7 +218,7 @@ func resolveStreamingDirectory( tracker.printProgress($0) } tracker.finish() - guard topLevel.allSatisfy({ isPresentFile($0, in: snapshot) }), + guard requiredListedFiles(files, in: snapshot).allSatisfy({ isPresentFile($0, in: snapshot) }), hasModelConfigAndTokenizer(in: snapshot) else { throw ModelDownloadIncomplete(modelId: modelId, directory: snapshot) diff --git a/tests/SwiftLMTests/StreamingDirectoryTests.swift b/tests/SwiftLMTests/StreamingDirectoryTests.swift index d555cfb..778a597 100644 --- a/tests/SwiftLMTests/StreamingDirectoryTests.swift +++ b/tests/SwiftLMTests/StreamingDirectoryTests.swift @@ -61,6 +61,41 @@ final class StreamingDirectoryTests: XCTestCase { XCTAssertEqual(localWeightState(in: dir), .incomplete) } + /// OptiQ VLMs index their vision tower in optiq/; a download of only model*.safetensors + /// has every top-level shard but not that file. + private func indexWithNestedVision() throws { + var map = Dictionary(uniqueKeysWithValues: (1 ... 2).map { + ("layer\($0).weight", String(format: "model-%05d-of-%05d.safetensors", $0, 2)) + }) + map["vision_tower.weight"] = "optiq/optiq_vision.safetensors" + let json = try JSONSerialization.data(withJSONObject: ["weight_map": map]) + try json.write(to: dir.appendingPathComponent("model.safetensors.index.json")) + } + + func testIndexedFileMissingFromSubdirectoryIsIncomplete() throws { + try indexWithNestedVision() + try shards([1, 2], of: 2) + XCTAssertEqual(localWeightState(in: dir), .incomplete) + } + + func testIndexedSubdirectoryFilePresentIsComplete() throws { + try indexWithNestedVision() + try shards([1, 2], of: 2) + try FileManager.default.createDirectory( + at: dir.appendingPathComponent("optiq"), withIntermediateDirectories: true) + try write("optiq/optiq_vision.safetensors", "weights") + XCTAssertEqual(localWeightState(in: dir), .complete) + } + + func testRequiredListedFilesIncludeIndexedNestedFiles() throws { + try indexWithNestedVision() + let listed = ["config.json", "model-00001-of-00002.safetensors", + "optiq/optiq_vision.safetensors", "optiq/mtp.safetensors"] + XCTAssertEqual( + Set(requiredListedFiles(listed, in: dir)), + ["config.json", "model-00001-of-00002.safetensors", "optiq/optiq_vision.safetensors"]) + } + func testSingleModelFileIsComplete() throws { try write("model.safetensors", "weights") XCTAssertEqual(localWeightState(in: dir), .complete)