diff --git a/cmd/jean/resource_windows_amd64.syso b/cmd/jean/resource_windows_amd64.syso index 9cdc0ea..94d5d81 100644 Binary files a/cmd/jean/resource_windows_amd64.syso and b/cmd/jean/resource_windows_amd64.syso differ diff --git a/cmd/jean/resource_windows_arm64.syso b/cmd/jean/resource_windows_arm64.syso index cb06bd5..c58ae6a 100644 Binary files a/cmd/jean/resource_windows_arm64.syso and b/cmd/jean/resource_windows_arm64.syso differ diff --git a/cmd/jean/versioninfo.json b/cmd/jean/versioninfo.json index 739d384..8d55758 100644 --- a/cmd/jean/versioninfo.json +++ b/cmd/jean/versioninfo.json @@ -3,13 +3,13 @@ "FileVersion": { "Major": 0, "Minor": 6, - "Patch": 4, + "Patch": 5, "Build": 0 }, "ProductVersion": { "Major": 0, "Minor": 6, - "Patch": 4, + "Patch": 5, "Build": 0 }, "FileFlagsMask": "3f", @@ -25,7 +25,7 @@ "LegalCopyright": "Copyright (c) 2026 AJEAN contributors. MIT License.", "OriginalFilename": "jean.exe", "ProductName": "AJEAN", - "ProductVersion": "0.6.4", + "ProductVersion": "0.6.5", "Comments": "https://github.com/nathaninline/jean — projet open source (MIT)" }, "VarFileInfo": { diff --git a/internal/jean/backend_models.go b/internal/jean/backend_models.go index b019b00..2bef28c 100644 --- a/internal/jean/backend_models.go +++ b/internal/jean/backend_models.go @@ -1,6 +1,7 @@ package jean import ( + "context" "encoding/json" "fmt" "io" @@ -10,8 +11,10 @@ import ( "path" "path/filepath" "regexp" + "strconv" "strings" "sync" + "sync/atomic" "time" ) @@ -156,9 +159,14 @@ type dlState struct { URL string `json:"url"` Total int64 `json:"total"` Done int64 `json:"done"` + Speed int64 `json:"speed"` // bytes/s, smoothed over the last samples + Conns int `json:"conns"` // parallel connections actually used Finished bool `json:"finished"` + Canceled bool `json:"canceled"` Err string `json:"error"` StartedAt int64 `json:"started_at"` + + cancel context.CancelFunc `json:"-"` // set while in flight, cleared on finish } var ( @@ -166,6 +174,44 @@ var ( dlDownloads = map[string]*dlState{} // keyed by filename ) +// dlClient is shared by all download workers so connections to the HF CDN are +// pooled and reused across chunks instead of re-handshaking TLS each time. +var dlClient = &http.Client{ + Timeout: 0, // large files: no overall timeout + Transport: &http.Transport{ + Proxy: http.ProxyFromEnvironment, + MaxIdleConns: 64, + MaxIdleConnsPerHost: 64, + MaxConnsPerHost: 0, + IdleConnTimeout: 90 * time.Second, + TLSHandshakeTimeout: 20 * time.Second, + ExpectContinueTimeout: 1 * time.Second, + // HTTP/1.1: truly parallel sockets, no shared h2 flow-control window. + ForceAttemptHTTP2: false, + WriteBufferSize: 64 << 10, + ReadBufferSize: 256 << 10, + }, +} + +// dlConns is the number of parallel range requests used per download. +// Overridable with JEAN_DL_CONNS (1 disables parallelism). +func dlConns() int { + n := 8 + if v := strings.TrimSpace(os.Getenv("JEAN_DL_CONNS")); v != "" { + if p, err := strconv.Atoi(v); err == nil && p > 0 { + n = p + } + } + if n > 16 { + n = 16 + } + return n +} + +// dlMinChunk is the smallest slice worth a dedicated connection (16 MiB), so a +// small file doesn't get split into a swarm of tiny requests. +const dlMinChunk = 16 << 20 + // normalizeHFURL turns a Hugging Face "blob" page URL into a direct "resolve" // download URL, and leaves already-direct URLs untouched. Returns the URL to // fetch and the target filename. @@ -228,79 +274,177 @@ func handleModelDownload(w http.ResponseWriter, r *http.Request) { sendJSON(w, 409, map[string]any{"ok": false, "error": "le modèle existe déjà: " + name}) return } - st := &dlState{Filename: name, URL: dlURL, StartedAt: time.Now().Unix()} + ctx, cancel := context.WithCancel(context.Background()) + st := &dlState{Filename: name, URL: dlURL, StartedAt: time.Now().Unix(), cancel: cancel} dlDownloads[name] = st dlMu.Unlock() - go runDownload(st, dlURL, dest) + go runDownload(ctx, st, dlURL, dest) sendJSON(w, 200, map[string]any{"ok": true, "filename": name}) } -// runDownload streams the URL to a .part file then renames it on success. -func runDownload(st *dlState, dlURL, dest string) { - finish := func(e error) { - dlMu.Lock() - if e != nil { - st.Err = e.Error() - } - st.Finished = true - dlMu.Unlock() - } - - req, err := http.NewRequest("GET", dlURL, nil) +// dlRequest builds a GET for the download URL, carrying the HF token when set +// (gated/private repos) and an optional Range header. +func dlRequest(ctx context.Context, dlURL, rng string) (*http.Request, error) { + req, err := http.NewRequestWithContext(ctx, "GET", dlURL, nil) if err != nil { - finish(err) - return + return nil, err } // HF gated/private repos may need a token; reuse the same key store if set. if k := os.Getenv("HF_TOKEN"); k != "" { req.Header.Set("Authorization", "Bearer "+k) } - client := &http.Client{Timeout: 0} // large files: no overall timeout - resp, err := client.Do(req) + req.Header.Set("User-Agent", "jean/"+Version) + req.Header.Set("Accept-Encoding", "identity") // never gzip a .gguf: it breaks ranges + if rng != "" { + req.Header.Set("Range", rng) + } + return req, nil +} + +// contentRangeTotal parses the total size out of a "bytes 0-0/12345" header. +func contentRangeTotal(v string) int64 { + i := strings.LastIndexByte(v, '/') + if i < 0 { + return 0 + } + n, err := strconv.ParseInt(strings.TrimSpace(v[i+1:]), 10, 64) + if err != nil || n <= 0 { + return 0 + } + return n +} + +// dlProbe asks the server for the first byte to learn the total size and +// whether ranges are supported (206 + Content-Range). +func dlProbe(ctx context.Context, dlURL string) (total int64, ranged bool, err error) { + req, err := dlRequest(ctx, dlURL, "bytes=0-0") + if err != nil { + return 0, false, err + } + resp, err := dlClient.Do(req) + if err != nil { + return 0, false, err + } + defer resp.Body.Close() + _, _ = io.Copy(io.Discard, resp.Body) + switch resp.StatusCode { + case 206: + if t := contentRangeTotal(resp.Header.Get("Content-Range")); t > 0 { + return t, true, nil + } + return 0, false, nil + case 200: + // Server ignored the Range: single stream, ContentLength is the size. + return resp.ContentLength, false, nil + default: + return 0, false, fmt.Errorf("HTTP %d depuis la source", resp.StatusCode) + } +} + +// runDownload fetches the URL into a .part file then renames it on success. +// When the source honours byte ranges (Hugging Face's CDN does), the file is +// split across several connections written in place with WriteAt, which is what +// makes a multi-GB .gguf saturate the link instead of a single TCP stream. +// Cancelling the context aborts every worker and the partial .part is removed, +// so a cancelled download leaves nothing behind. +func runDownload(ctx context.Context, st *dlState, dlURL, dest string) { + tmp := dest + ".part" + var done int64 // atomic: total bytes written across all workers + + finish := func(e error) { + dlMu.Lock() + switch { + case ctx.Err() != nil: + st.Canceled = true + st.Speed = 0 + case e != nil: + st.Err = e.Error() + st.Speed = 0 + default: + st.Done = atomic.LoadInt64(&done) + } + st.Finished = true + st.cancel = nil + dlMu.Unlock() + } + + total, ranged, err := dlProbe(ctx, dlURL) if err != nil { finish(err) return } - defer resp.Body.Close() - if resp.StatusCode != 200 { - finish(fmt.Errorf("HTTP %d depuis la source", resp.StatusCode)) - return + + conns := 1 + if ranged && total > 0 { + conns = dlConns() + if max := int((total + dlMinChunk - 1) / dlMinChunk); conns > max { + conns = max + } + if conns < 1 { + conns = 1 + } } + dlMu.Lock() - st.Total = resp.ContentLength + st.Total = total + st.Conns = conns dlMu.Unlock() - tmp := dest + ".part" f, err := os.Create(tmp) if err != nil { finish(err) return } - buf := make([]byte, 1<<20) // 1 MiB - for { - n, rerr := resp.Body.Read(buf) - if n > 0 { - if _, werr := f.Write(buf[:n]); werr != nil { - f.Close() - _ = os.Remove(tmp) - finish(werr) - return - } - dlMu.Lock() - st.Done += int64(n) - dlMu.Unlock() - } - if rerr == io.EOF { - break - } - if rerr != nil { - f.Close() - _ = os.Remove(tmp) - finish(rerr) + fail := func(e error) { + f.Close() + _ = os.Remove(tmp) + finish(e) + } + if conns > 1 { + // Preallocate so the filesystem can lay the file out contiguously and + // concurrent WriteAt calls never race to extend it. + if err := f.Truncate(total); err != nil { + fail(err) return } } + + // Publish progress + a smoothed speed once a second. + stop := make(chan struct{}) + go func() { + t := time.NewTicker(time.Second) + defer t.Stop() + last, lastAt := int64(0), time.Now() + for { + select { + case <-stop: + return + case now := <-t.C: + cur := atomic.LoadInt64(&done) + dt := now.Sub(lastAt).Seconds() + dlMu.Lock() + st.Done = cur + if dt > 0 { + inst := int64(float64(cur-last) / dt) + if st.Speed == 0 { + st.Speed = inst + } else { + st.Speed = (st.Speed*2 + inst) / 3 // EMA, smooths CDN jitter + } + } + dlMu.Unlock() + last, lastAt = cur, now + } + } + }() + + err = dlFetch(ctx, f, dlURL, total, conns, &done) + close(stop) + if err != nil { + fail(err) // includes cancellation: the .part is deleted either way + return + } if err := f.Close(); err != nil { _ = os.Remove(tmp) finish(err) @@ -314,6 +458,159 @@ func runDownload(st *dlState, dlURL, dest string) { finish(nil) } +// dlFetch writes the whole body into f, either as one stream or as `conns` +// parallel byte ranges. done is incremented atomically as bytes land on disk. +func dlFetch(ctx context.Context, f *os.File, dlURL string, total int64, conns int, done *int64) error { + if conns <= 1 { + return dlChunk(ctx, f, dlURL, 0, total-1, total <= 0, done) + } + size := total / int64(conns) + var wg sync.WaitGroup + errs := make([]error, conns) + for i := 0; i < conns; i++ { + start := int64(i) * size + end := start + size - 1 + if i == conns-1 { + end = total - 1 + } + wg.Add(1) + go func(i int, start, end int64) { + defer wg.Done() + errs[i] = dlChunk(ctx, f, dlURL, start, end, false, done) + }(i, start, end) + } + wg.Wait() + for _, e := range errs { + if e != nil { + return e + } + } + return nil +} + +// dlChunk downloads [start,end] into f at the right offset, retrying from where +// it stopped if the connection drops mid-chunk. With whole=true it streams the +// entire body sequentially (server without range support, unknown size). +func dlChunk(ctx context.Context, f *os.File, dlURL string, start, end int64, whole bool, done *int64) error { + const attempts = 4 + pos := start + var lastErr error + for try := 0; try < attempts; try++ { + if err := ctx.Err(); err != nil { + return err // cancelled: never retry + } + if try > 0 { + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(time.Duration(try) * time.Second): + } + } + rng := "" + if !whole { + if pos > end { + return nil + } + rng = fmt.Sprintf("bytes=%d-%d", pos, end) + } else if pos > start { + rng = fmt.Sprintf("bytes=%d-", pos) // best-effort resume + } + req, err := dlRequest(ctx, dlURL, rng) + if err != nil { + return err + } + resp, err := dlClient.Do(req) + if err != nil { + lastErr = err + continue + } + if resp.StatusCode != 200 && resp.StatusCode != 206 { + resp.Body.Close() + return fmt.Errorf("HTTP %d depuis la source", resp.StatusCode) + } + if resp.StatusCode == 200 && pos > start { + // Resume refused: the body restarts from 0, rewind our bookkeeping. + atomic.AddInt64(done, start-pos) + pos = start + } + n, cerr := dlCopy(f, resp.Body, pos, done) + resp.Body.Close() + pos += n + if cerr == nil { + return nil + } + lastErr = cerr + } + return lastErr +} + +// dlCopy streams src into f starting at off, reporting bytes written. It +// returns the byte count even on error so the caller can resume. +func dlCopy(f *os.File, src io.Reader, off int64, done *int64) (int64, error) { + buf := make([]byte, 1<<20) // 1 MiB + var written int64 + for { + n, rerr := src.Read(buf) + if n > 0 { + if _, werr := f.WriteAt(buf[:n], off+written); werr != nil { + return written, werr + } + written += int64(n) + atomic.AddInt64(done, int64(n)) + } + if rerr == io.EOF { + return written, nil + } + if rerr != nil { + return written, rerr + } + } +} + +// handleModelDownloadCancel aborts an in-flight download; runDownload then +// deletes its .part file, so nothing partial survives. +func handleModelDownloadCancel(w http.ResponseWriter, r *http.Request) { + var req struct { + Filename string `json:"filename"` + } + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + sendJSON(w, 400, map[string]any{"ok": false, "error": err.Error()}) + return + } + name := filepath.Base(strings.TrimSpace(req.Filename)) + dlMu.Lock() + st, ok := dlDownloads[name] + var cancel context.CancelFunc + if ok { + cancel = st.cancel + } + dlMu.Unlock() + if !ok { + sendJSON(w, 404, map[string]any{"ok": false, "error": "aucun téléchargement pour " + name}) + return + } + if cancel != nil { + cancel() + } + sendJSON(w, 200, map[string]any{"ok": true}) +} + +// cleanStalePartFiles removes leftover *.gguf.part files in JEAN_HOME at +// startup. A download killed by a crash or a service restart can't be resumed +// (its state lived in memory), so the partial file would otherwise sit there +// forever eating disk. +func cleanStalePartFiles() { + matches, err := filepath.Glob(filepath.Join(JeanHome(), "*.gguf.part")) + if err != nil { + return + } + for _, p := range matches { + if err := os.Remove(p); err == nil { + fmt.Printf("[models] téléchargement incomplet supprimé : %s\n", filepath.Base(p)) + } + } +} + // handleModelDownloadStatus returns the state of all known downloads this run. func handleModelDownloadStatus(w http.ResponseWriter, r *http.Request) { dlMu.Lock() diff --git a/internal/jean/backend_models_dl_test.go b/internal/jean/backend_models_dl_test.go new file mode 100644 index 0000000..b295072 --- /dev/null +++ b/internal/jean/backend_models_dl_test.go @@ -0,0 +1,134 @@ +package jean + +import ( + "bytes" + "context" + "math/rand" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strconv" + "testing" + "time" +) + +// serveBlob serves data with (or without) byte-range support. +func serveBlob(data []byte, ranges bool) *httptest.Server { + return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if !ranges { + w.Header().Set("Content-Length", strconv.Itoa(len(data))) + _, _ = w.Write(data) + return + } + http.ServeContent(w, r, "m.gguf", time.Time{}, bytes.NewReader(data)) + })) +} + +// Un téléchargement annulé ne doit laisser NI le .gguf final NI le .part. +func TestRunDownloadCancelLeavesNothing(t *testing.T) { + t.Setenv("JEAN_DL_CONNS", "4") + release := make(chan struct{}) + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Header.Get("Range") == "bytes=0-0" { // sonde : réponse immédiate + w.Header().Set("Content-Range", "bytes 0-0/"+strconv.Itoa(64<<20)) + w.WriteHeader(206) + _, _ = w.Write([]byte{0}) + return + } + // Corps qui traîne : l'annulation doit l'interrompre. + w.WriteHeader(206) + w.(http.Flusher).Flush() + select { + case <-release: + case <-r.Context().Done(): + } + })) + defer srv.Close() + defer close(release) + + dir := t.TempDir() + dest := filepath.Join(dir, "m.gguf") + ctx, cancel := context.WithCancel(context.Background()) + st := &dlState{Filename: "m.gguf", cancel: cancel} + done := make(chan struct{}) + go func() { runDownload(ctx, st, srv.URL+"/m.gguf", dest); close(done) }() + + time.Sleep(300 * time.Millisecond) + cancel() + select { + case <-done: + case <-time.After(10 * time.Second): + t.Fatal("runDownload n'a pas rendu la main après annulation") + } + + if !st.Canceled || !st.Finished { + t.Fatalf("état attendu annulé+terminé, got canceled=%v finished=%v err=%q", st.Canceled, st.Finished, st.Err) + } + if _, err := os.Stat(dest + ".part"); !os.IsNotExist(err) { + t.Fatal(".part laissé sur le disque après annulation") + } + if _, err := os.Stat(dest); !os.IsNotExist(err) { + t.Fatal("fichier final créé alors que le téléchargement a été annulé") + } +} + +// cleanStalePartFiles doit balayer les .part orphelins sans toucher aux .gguf. +func TestCleanStalePartFiles(t *testing.T) { + dir := t.TempDir() + t.Setenv("JEAN_HOME", dir) + keep := filepath.Join(dir, "bon.gguf") + stale := filepath.Join(dir, "coupe.gguf.part") + for _, p := range []string{keep, stale} { + if err := os.WriteFile(p, []byte("x"), 0o644); err != nil { + t.Fatal(err) + } + } + cleanStalePartFiles() + if _, err := os.Stat(stale); !os.IsNotExist(err) { + t.Fatal(".part orphelin non supprimé") + } + if _, err := os.Stat(keep); err != nil { + t.Fatal(".gguf valide supprimé par erreur") + } +} + +func TestRunDownloadParallelAndFallback(t *testing.T) { + data := make([]byte, 48<<20) // > 3×dlMinChunk so the split actually kicks in + rand.New(rand.NewSource(1)).Read(data) + t.Setenv("JEAN_DL_CONNS", "4") + + for _, ranges := range []bool{true, false} { + srv := serveBlob(data, ranges) + dir := t.TempDir() + dest := filepath.Join(dir, "m.gguf") + st := &dlState{Filename: "m.gguf"} + runDownload(context.Background(), st, srv.URL+"/m.gguf", dest) + srv.Close() + + if st.Err != "" { + t.Fatalf("ranges=%v: erreur %s", ranges, st.Err) + } + got, err := os.ReadFile(dest) + if err != nil { + t.Fatalf("ranges=%v: %v", ranges, err) + } + if len(got) != len(data) { + t.Fatalf("ranges=%v: taille %d != %d", ranges, len(got), len(data)) + } + for i := range got { + if got[i] != data[i] { + t.Fatalf("ranges=%v: octet %d différent", ranges, i) + } + } + if st.Done != int64(len(data)) { + t.Fatalf("ranges=%v: done=%d", ranges, st.Done) + } + if ranges && st.Conns < 2 { + t.Fatalf("attendu du parallélisme, conns=%d", st.Conns) + } + if _, err := os.Stat(dest + ".part"); !os.IsNotExist(err) { + t.Fatalf("ranges=%v: .part laissé derrière", ranges) + } + } +} diff --git a/internal/jean/run.go b/internal/jean/run.go index 8b75ca0..22a48dc 100644 --- a/internal/jean/run.go +++ b/internal/jean/run.go @@ -10,7 +10,7 @@ import ( "strings" ) -const Version = "0.6.4" +const Version = "0.6.5" // Main est le vrai main() du binaire (cmd/jean ne fait que l'appeler). func Main() { diff --git a/internal/jean/ui/index.html b/internal/jean/ui/index.html index c3c9303..c2cc66e 100644 --- a/internal/jean/ui/index.html +++ b/internal/jean/ui/index.html @@ -560,6 +560,12 @@ textarea.sctl{min-height:76px;resize:vertical;line-height:1.5} .pe-selc select{appearance:none;-webkit-appearance:none;border:none;background:transparent;color:var(--accent);font:inherit;text-align:right;text-align-last:right;padding:4px 16px 4px 0;cursor:pointer;max-width:100%;text-overflow:ellipsis} .pe-selc select:focus{outline:none} .pe-selc::after{content:"›";position:absolute;right:0;transform:rotate(90deg);color:var(--dim);font-size:14px;pointer-events:none} +/* Variante pleine largeur : la liste native s'ouvre alignée à gauche de la + ligne (et non collée au bord droit) et le nom de fichier a la place de + s'afficher en entier. */ +.pe-selc.wide{width:100%} +.pe-selc.wide select{flex:1;text-align:left;text-align-last:left;padding-right:18px} +.pe-selc.wide::after{right:2px} /* Ligne « champ pleine largeur » (Nom, chemin, config brute…). */ .pe-row.stack{flex-direction:column;align-items:stretch;gap:8px} .pe-inp{width:100%;box-sizing:border-box;background:var(--bg);color:var(--text);border:1px solid var(--border);border-radius:8px;padding:9px 11px;font:inherit;transition:border-color .12s,box-shadow .12s} @@ -585,6 +591,17 @@ textarea.sctl{min-height:76px;resize:vertical;line-height:1.5} .pe-seg .be-note{display:none} .pe-seg-note{font-size:11px;color:var(--dim);padding:0 2px;min-height:14px} .pe-note{font-size:11px;color:var(--dim);line-height:1.45} +/* Barre de progression (téléchargement de modèle). Sans total connu, elle + défile en boucle via .indet. */ +.pe-bar{height:4px;border-radius:999px;background:color-mix(in srgb,var(--dim) 26%,transparent);overflow:hidden} +.pe-bar>i{display:block;height:100%;width:0;border-radius:999px;background:var(--accent);transition:width .35s linear} +.pe-bar.indet>i{width:35%;animation:pe-bar-slide 1.2s ease-in-out infinite;transition:none} +.pe-bar.done>i{width:100%} +@keyframes pe-bar-slide{0%{margin-left:-35%}100%{margin-left:100%}} +/* Action discrète en texte (annuler un téléchargement…). */ +.pe-link{align-self:flex-start;border:none;background:none;padding:0;font:inherit;font-size:11px;color:var(--dim);text-decoration:underline;cursor:pointer} +.pe-link:hover{color:var(--text)} +.pe-link:disabled{opacity:.5;cursor:default;text-decoration:none} /* Interrupteur à glissière. */ .pe-switch{position:relative;flex:none;width:40px;height:24px} .pe-switch input{position:absolute;opacity:0;width:0;height:0} @@ -1189,9 +1206,9 @@ button:hover{border-color:var(--dim);color:var(--text)}