mirror of
https://github.com/R0m1k3/Loki.git
synced 2026-10-11 17:26:57 +02:00
Envois simultanes serialises, et pas de controle d'espace disque
Deux defauts trouves en relisant le code de l'envoi en morceaux. Le verrou global etait tenu pendant l'ecriture disque et le renommage. Or le client lance un envoi PAR FICHIER JOINT, de front : trois fichiers s'attendaient donc les uns les autres, chacun bloquant les suivants jusqu'a sa derniere tranche. Le verrou global ne protege plus que la table des sessions ; l'ecriture se fait sous le verrou de la session. Un drapeau `busy` empeche le menage de fermer un fichier sous les pieds de la requete qui l'ecrit. Test : 3 fichiers de 20 Mo en parallele, contenus distincts, aucun melange. Rien ne verifiait la place disque : un envoi d'un gigaoctet pouvait remplir le volume de la machine qui fait tourner le modele, ou vivent aussi llama-server, la base et les journaux. Le client annonce la taille au premier morceau, le serveur refuse en 507 s'il ne reste pas la place plus 512 Mo de marge. La taille annoncee n'engage que le client : le plafond reel reste verifie morceau par morceau. Aussi : Sync() avant publication (un tampon non vide donne un fichier incomplet), et cote client une tranche vide sans marque de fin arrete le telechargement au lieu de boucler indefiniment.
This commit is contained in:
1 parent
78423276e7
commit
6d36464971
5 files changed
+146
-18
No files matched your search
+17
-3
@@ -1,4 +1,4 @@
|
|||||||
La limite d'envoi de fichiers passe de 24 Mo à 1 Go.
|
La limite d'envoi de fichiers passe de 24 Mo à 1 Go, et le téléchargement depuis ajean.link est réparé.
|
||||||
|
|
||||||
## Des fichiers d'un gigaoctet
|
## Des fichiers d'un gigaoctet
|
||||||
|
|
||||||
@@ -8,11 +8,23 @@ Le fichier part maintenant par tranches de 8 Mo. Le navigateur n'en lit qu'une
|
|||||||
|
|
||||||
Au passage, la vignette affiche l'avancement en pourcentage : sur un gros fichier, l'anneau qui tournait sans rien dire n'était pas d'un grand secours.
|
Au passage, la vignette affiche l'avancement en pourcentage : sur un gros fichier, l'anneau qui tournait sans rien dire n'était pas d'un grand secours.
|
||||||
|
|
||||||
|
Avant d'entamer un envoi, AJEAN vérifie qu'il reste de la place sur le disque et refuse franchement s'il n'y en a pas — remplir le volume de la machine qui fait tourner le modèle serait autrement plus ennuyeux qu'un envoi rejeté.
|
||||||
|
|
||||||
Un envoi interrompu — onglet fermé, réseau coupé, service redémarré — ne laisse rien derrière lui : le fichier partiel est effacé au bout de dix minutes, et au démarrage du service.
|
Un envoi interrompu — onglet fermé, réseau coupé, service redémarré — ne laisse rien derrière lui : le fichier partiel est effacé au bout de dix minutes, et au démarrage du service.
|
||||||
|
|
||||||
Ce qui n'a pas changé : le dépôt n'a toujours lieu qu'à l'envoi du message, et l'accès distant fonctionne pareil, les tranches passant par le tunnel chiffré comme le reste.
|
## Télécharger un fichier depuis ajean.link
|
||||||
|
|
||||||
Pensez tout de même à la place disque de la machine qui héberge AJEAN : les fichiers envoyés s'accumulent dans `uploads/`, à l'intérieur du dossier de travail de l'IA.
|
Depuis l'accès distant, cliquer sur un fichier proposé par l'IA téléchargeait un contenu JSON illisible à la place du fichier.
|
||||||
|
|
||||||
|
La cause tient à la façon dont l'accès distant chiffre les échanges : toute réponse est réemballée en JSON avant de vous parvenir, et des données binaires n'y survivent pas. Le fichier était donc détruit en chemin, quel que soit son type.
|
||||||
|
|
||||||
|
Votre serveur signale maintenant à l'interface qu'elle est derrière le tunnel, et lui envoie le fichier sous une forme qui traverse intact, toujours en tranches pour ne pas le charger entier en mémoire. En local, rien ne change : le fichier est servi directement, sans détour.
|
||||||
|
|
||||||
|
Le même défaut abîmait l'export de conversation au format Markdown, qui revenait truffé de guillemets et d'échappements. Corrigé aussi.
|
||||||
|
|
||||||
|
## Aussi
|
||||||
|
|
||||||
|
- Plusieurs fichiers joints en même temps s'envoient désormais en parallèle. Ils s'attendaient les uns les autres, chaque envoi bloquant les suivants jusqu'à sa dernière tranche.
|
||||||
|
|
||||||
## Mise à jour
|
## Mise à jour
|
||||||
|
|
||||||
@@ -25,3 +37,5 @@ Puis redémarrez l'interface :
|
|||||||
```bash
|
```bash
|
||||||
ajean ui restart
|
ajean ui restart
|
||||||
```
|
```
|
||||||
|
|
||||||
|
Pensez à la place disque de la machine qui héberge AJEAN : les fichiers envoyés s'accumulent dans `uploads/`, à l'intérieur du dossier de travail de l'IA, et rien ne les efface pour l'instant.
|
||||||
@@ -5114,8 +5114,12 @@ async function fetchFileB64(path, size, chip){
|
|||||||
const r=await jfetch('/api/chat/file?b64=1&len='+DL_CHUNK+'&offset='+off+'&path='+encodeURIComponent(path));
|
const r=await jfetch('/api/chat/file?b64=1&len='+DL_CHUNK+'&offset='+off+'&path='+encodeURIComponent(path));
|
||||||
const j=await r.json();
|
const j=await r.json();
|
||||||
if(!r.ok || !j.ok) throw new Error(j.error||('HTTP '+r.status));
|
if(!r.ok || !j.ok) throw new Error(j.error||('HTTP '+r.status));
|
||||||
parts.push(b64ToBytes(j.data||''));
|
const bytes=b64ToBytes(j.data||'');
|
||||||
off=(j.offset||0)+(parts[parts.length-1].length);
|
parts.push(bytes);
|
||||||
|
// Une tranche vide sans marque de fin ne ferait pas avancer la boucle : on
|
||||||
|
// s'arrête plutôt que de tourner indéfiniment sur le même octet.
|
||||||
|
if(!bytes.length && !j.eof) throw new Error('transfert interrompu');
|
||||||
|
off=(j.offset||0)+bytes.length;
|
||||||
if(chip && size) chip.title='téléchargement '+Math.round(off*100/size)+' %';
|
if(chip && size) chip.title='téléchargement '+Math.round(off*100/size)+' %';
|
||||||
if(j.eof) break;
|
if(j.eof) break;
|
||||||
} while(off<size);
|
} while(off<size);
|
||||||
@@ -5264,7 +5268,9 @@ async function uploadAttach(a){
|
|||||||
const end=Math.min(off+ATTACH_CHUNK, a.size);
|
const end=Math.min(off+ATTACH_CHUNK, a.size);
|
||||||
const data=await readAsBase64(a.file.slice(off, end));
|
const data=await readAsBase64(a.file.slice(off, end));
|
||||||
const last=end>=a.size;
|
const last=end>=a.size;
|
||||||
const j=await sendChunk({name:a.name, data:data, id:id, more:!last});
|
// `size` au premier morceau : le serveur vérifie la place disque AVANT
|
||||||
|
// d'entamer un envoi d'un gigaoctet, plutôt que de le découvrir à la fin.
|
||||||
|
const j=await sendChunk({name:a.name, data:data, id:id, more:!last, size:(id?0:a.size)});
|
||||||
if(j.id) id=j.id;
|
if(j.id) id=j.id;
|
||||||
off=end;
|
off=end;
|
||||||
a.sent=off; renderAttach();
|
a.sent=off; renderAttach();
|
||||||
|
|||||||
@@ -93,8 +93,12 @@ async function fetchFileB64(path, size, chip){
|
|||||||
const r=await jfetch('/api/chat/file?b64=1&len='+DL_CHUNK+'&offset='+off+'&path='+encodeURIComponent(path));
|
const r=await jfetch('/api/chat/file?b64=1&len='+DL_CHUNK+'&offset='+off+'&path='+encodeURIComponent(path));
|
||||||
const j=await r.json();
|
const j=await r.json();
|
||||||
if(!r.ok || !j.ok) throw new Error(j.error||('HTTP '+r.status));
|
if(!r.ok || !j.ok) throw new Error(j.error||('HTTP '+r.status));
|
||||||
parts.push(b64ToBytes(j.data||''));
|
const bytes=b64ToBytes(j.data||'');
|
||||||
off=(j.offset||0)+(parts[parts.length-1].length);
|
parts.push(bytes);
|
||||||
|
// Une tranche vide sans marque de fin ne ferait pas avancer la boucle : on
|
||||||
|
// s'arrête plutôt que de tourner indéfiniment sur le même octet.
|
||||||
|
if(!bytes.length && !j.eof) throw new Error('transfert interrompu');
|
||||||
|
off=(j.offset||0)+bytes.length;
|
||||||
if(chip && size) chip.title='téléchargement '+Math.round(off*100/size)+' %';
|
if(chip && size) chip.title='téléchargement '+Math.round(off*100/size)+' %';
|
||||||
if(j.eof) break;
|
if(j.eof) break;
|
||||||
} while(off<size);
|
} while(off<size);
|
||||||
@@ -243,7 +247,9 @@ async function uploadAttach(a){
|
|||||||
const end=Math.min(off+ATTACH_CHUNK, a.size);
|
const end=Math.min(off+ATTACH_CHUNK, a.size);
|
||||||
const data=await readAsBase64(a.file.slice(off, end));
|
const data=await readAsBase64(a.file.slice(off, end));
|
||||||
const last=end>=a.size;
|
const last=end>=a.size;
|
||||||
const j=await sendChunk({name:a.name, data:data, id:id, more:!last});
|
// `size` au premier morceau : le serveur vérifie la place disque AVANT
|
||||||
|
// d'entamer un envoi d'un gigaoctet, plutôt que de le découvrir à la fin.
|
||||||
|
const j=await sendChunk({name:a.name, data:data, id:id, more:!last, size:(id?0:a.size)});
|
||||||
if(j.id) id=j.id;
|
if(j.id) id=j.id;
|
||||||
off=end;
|
off=end;
|
||||||
a.sent=off; renderAttach();
|
a.sent=off; renderAttach();
|
||||||
|
|||||||
@@ -278,6 +278,10 @@ type uploadReq struct {
|
|||||||
// le fichier. Un envoi en un seul morceau (More absent) reste valable.
|
// le fichier. Un envoi en un seul morceau (More absent) reste valable.
|
||||||
ID string `json:"id"`
|
ID string `json:"id"`
|
||||||
More bool `json:"more"`
|
More bool `json:"more"`
|
||||||
|
// Size = taille totale annoncée au PREMIER morceau, pour vérifier l'espace
|
||||||
|
// disque avant d'entamer un envoi d'un gigaoctet. Purement indicatif : le
|
||||||
|
// vrai plafond reste vérifié morceau par morceau.
|
||||||
|
Size int64 `json:"size"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// upSession = un envoi en cours, adossé à un fichier .part sur le disque.
|
// upSession = un envoi en cours, adossé à un fichier .part sur le disque.
|
||||||
@@ -288,6 +292,12 @@ type uploadReq struct {
|
|||||||
// base64 reçu et le binaire décodé), ce qui plafonnait l'envoi à quelques Mo
|
// base64 reçu et le binaire décodé), ce qui plafonnait l'envoi à quelques Mo
|
||||||
// utilisables sans faire gonfler le process.
|
// utilisables sans faire gonfler le process.
|
||||||
type upSession struct {
|
type upSession struct {
|
||||||
|
// mu sérialise les écritures d'UNE session. Le verrou global ne couvre que la
|
||||||
|
// table : le tenir pendant l'écriture disque mettait tous les envois à la
|
||||||
|
// queue leu leu, alors que le client en lance plusieurs de front (un par
|
||||||
|
// fichier joint).
|
||||||
|
mu sync.Mutex
|
||||||
|
busy bool // écriture en cours : le ménage doit passer son tour
|
||||||
f *os.File
|
f *os.File
|
||||||
name string
|
name string
|
||||||
written int64
|
written int64
|
||||||
@@ -299,6 +309,11 @@ var (
|
|||||||
upSessions = map[string]*upSession{}
|
upSessions = map[string]*upSession{}
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// uploadSpaceMargin : place qu'on refuse d'entamer sur le disque. Remplir le
|
||||||
|
// volume de la machine qui fait tourner le modèle est autrement plus grave que
|
||||||
|
// de refuser un envoi — llama-server, la base et les journaux vivent dessus.
|
||||||
|
const uploadSpaceMargin = 512 << 20
|
||||||
|
|
||||||
// upSessionTTL : au-delà, un envoi interrompu (onglet fermé, réseau coupé) est
|
// upSessionTTL : au-delà, un envoi interrompu (onglet fermé, réseau coupé) est
|
||||||
// abandonné et son .part supprimé. Sans ça, un fichier à moitié transféré
|
// abandonné et son .part supprimé. Sans ça, un fichier à moitié transféré
|
||||||
// resterait ouvert et occuperait le disque indéfiniment.
|
// resterait ouvert et occuperait le disque indéfiniment.
|
||||||
@@ -309,7 +324,9 @@ const upSessionTTL = 10 * time.Minute
|
|||||||
func sweepUploadSessions() {
|
func sweepUploadSessions() {
|
||||||
now := time.Now()
|
now := time.Now()
|
||||||
for id, s := range upSessions {
|
for id, s := range upSessions {
|
||||||
if now.Sub(s.last) > upSessionTTL {
|
// `busy` : une écriture est en cours dans cette session. La balayer
|
||||||
|
// fermerait le fichier sous les pieds de la requête qui l'écrit.
|
||||||
|
if !s.busy && now.Sub(s.last) > upSessionTTL {
|
||||||
name := s.f.Name()
|
name := s.f.Name()
|
||||||
s.f.Close()
|
s.f.Close()
|
||||||
os.Remove(name)
|
os.Remove(name)
|
||||||
@@ -343,25 +360,38 @@ func handleChatUpload(w http.ResponseWriter, r *http.Request) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Le verrou global ne protège QUE la table des sessions. L'écriture, elle, se
|
||||||
|
// fait sous le verrou de la session — sinon deux fichiers envoyés en même
|
||||||
|
// temps (le client les lance de front) s'attendraient l'un l'autre.
|
||||||
upMu.Lock()
|
upMu.Lock()
|
||||||
defer upMu.Unlock()
|
|
||||||
sweepUploadSessions()
|
sweepUploadSessions()
|
||||||
|
|
||||||
s := upSessions[body.ID]
|
s := upSessions[body.ID]
|
||||||
if s == nil {
|
if s == nil {
|
||||||
if body.ID != "" {
|
if body.ID != "" {
|
||||||
// L'envoi a expiré ou le serveur a redémarré en cours de route : le dire,
|
// L'envoi a expiré ou le serveur a redémarré en cours de route : le dire,
|
||||||
// plutôt que de recommencer un fichier à partir de son milieu.
|
// plutôt que de recommencer un fichier à partir de son milieu.
|
||||||
|
upMu.Unlock()
|
||||||
sendJSON(w, 409, map[string]any{"ok": false, "error": "envoi expiré — recommence le fichier"})
|
sendJSON(w, 409, map[string]any{"ok": false, "error": "envoi expiré — recommence le fichier"})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
dir, err := uploadsDir()
|
dir, err := uploadsDir()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
upMu.Unlock()
|
||||||
sendJSON(w, 500, map[string]any{"ok": false, "error": err.Error()})
|
sendJSON(w, 500, map[string]any{"ok": false, "error": err.Error()})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
// Espace disque : refuser franchement vaut mieux que remplir le volume de
|
||||||
|
// la machine qui fait tourner le modèle. La taille annoncée par le client
|
||||||
|
// n'engage que lui — le plafond réel reste vérifié morceau par morceau.
|
||||||
|
if free := diskFree(dir); free > 0 && body.Size > 0 && free < body.Size+uploadSpaceMargin {
|
||||||
|
upMu.Unlock()
|
||||||
|
sendJSON(w, 507, map[string]any{"ok": false,
|
||||||
|
"error": fmt.Sprintf("espace insuffisant : %s libres, %s nécessaires", humanBytes(free), humanBytes(body.Size+uploadSpaceMargin))})
|
||||||
|
return
|
||||||
|
}
|
||||||
f, err := os.CreateTemp(dir, ".upload-*.part")
|
f, err := os.CreateTemp(dir, ".upload-*.part")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
upMu.Unlock()
|
||||||
sendJSON(w, 500, map[string]any{"ok": false, "error": err.Error()})
|
sendJSON(w, 500, map[string]any{"ok": false, "error": err.Error()})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -370,12 +400,28 @@ func handleChatUpload(w http.ResponseWriter, r *http.Request) {
|
|||||||
upSessions[body.ID] = s
|
upSessions[body.ID] = s
|
||||||
}
|
}
|
||||||
s.last = time.Now()
|
s.last = time.Now()
|
||||||
|
s.busy = true
|
||||||
|
upMu.Unlock()
|
||||||
|
|
||||||
abort := func(code int, msg string) {
|
s.mu.Lock()
|
||||||
|
defer func() {
|
||||||
|
s.mu.Unlock()
|
||||||
|
upMu.Lock()
|
||||||
|
s.busy = false
|
||||||
|
upMu.Unlock()
|
||||||
|
}()
|
||||||
|
|
||||||
|
// drop ferme et oublie la session : sur erreur, et à la fin de l'envoi.
|
||||||
|
drop := func() string {
|
||||||
name := s.f.Name()
|
name := s.f.Name()
|
||||||
s.f.Close()
|
s.f.Close()
|
||||||
os.Remove(name)
|
upMu.Lock()
|
||||||
delete(upSessions, body.ID)
|
delete(upSessions, body.ID)
|
||||||
|
upMu.Unlock()
|
||||||
|
return name
|
||||||
|
}
|
||||||
|
abort := func(code int, msg string) {
|
||||||
|
os.Remove(drop())
|
||||||
sendJSON(w, code, map[string]any{"ok": false, "error": msg})
|
sendJSON(w, code, map[string]any{"ok": false, "error": msg})
|
||||||
}
|
}
|
||||||
if s.written+int64(len(raw)) > uploadMaxBytes {
|
if s.written+int64(len(raw)) > uploadMaxBytes {
|
||||||
@@ -396,12 +442,15 @@ func handleChatUpload(w http.ResponseWriter, r *http.Request) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// Dernier morceau : on ferme et on donne au fichier son vrai nom.
|
// Dernier morceau : on ferme et on donne au fichier son vrai nom. Le Close est
|
||||||
delete(upSessions, body.ID)
|
// dans drop() ; son erreur éventuelle est celle d'un tampon non vidé, donc
|
||||||
|
// d'un fichier incomplet — on ne le publie pas dans ce cas.
|
||||||
part := s.f.Name()
|
part := s.f.Name()
|
||||||
if err := s.f.Close(); err != nil {
|
syncErr := s.f.Sync()
|
||||||
|
drop()
|
||||||
|
if syncErr != nil {
|
||||||
os.Remove(part)
|
os.Remove(part)
|
||||||
sendJSON(w, 500, map[string]any{"ok": false, "error": err.Error()})
|
sendJSON(w, 500, map[string]any{"ok": false, "error": syncErr.Error()})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if s.written == 0 {
|
if s.written == 0 {
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import (
|
|||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -302,3 +303,55 @@ func TestChatFileB64ForE2E(t *testing.T) {
|
|||||||
t.Fatalf("fichier reconstitué : %d octets, attendu %d", len(got), len(raw))
|
t.Fatalf("fichier reconstitué : %d octets, attendu %d", len(got), len(raw))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Plusieurs fichiers partent DE FRONT depuis le navigateur (un envoi par pièce
|
||||||
|
// jointe). Les sessions doivent donc être indépendantes : le verrou global ne
|
||||||
|
// protège que la table, l'écriture se fait sous le verrou de la session.
|
||||||
|
func TestChatUploadConcurrentSessions(t *testing.T) {
|
||||||
|
const files, chunks = 4, 6
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
paths := make([]string, files)
|
||||||
|
errs := make([]string, files)
|
||||||
|
for i := 0; i < files; i++ {
|
||||||
|
wg.Add(1)
|
||||||
|
go func(i int) {
|
||||||
|
defer wg.Done()
|
||||||
|
id, want := "", ""
|
||||||
|
for c := 0; c < chunks; c++ {
|
||||||
|
piece := strings.Repeat(string(rune('A'+i)), 100+c)
|
||||||
|
want += piece
|
||||||
|
code, resp := uploadChunk(t, map[string]any{
|
||||||
|
"name": fmt.Sprintf("concurrent-%d.txt", i),
|
||||||
|
"data": base64.StdEncoding.EncodeToString([]byte(piece)),
|
||||||
|
"id": id, "more": c < chunks-1,
|
||||||
|
})
|
||||||
|
if code != 200 {
|
||||||
|
errs[i] = fmt.Sprintf("morceau %d : code %d (%v)", c, code, resp["error"])
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if v, ok := resp["id"].(string); ok && v != "" {
|
||||||
|
id = v
|
||||||
|
}
|
||||||
|
if c == chunks-1 {
|
||||||
|
abs, _ := resp["abs"].(string)
|
||||||
|
paths[i] = abs
|
||||||
|
got, err := os.ReadFile(abs)
|
||||||
|
if err != nil || string(got) != want {
|
||||||
|
errs[i] = fmt.Sprintf("contenu melange : %d octets lus, %d attendus (err %v)", len(got), len(want), err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}(i)
|
||||||
|
}
|
||||||
|
wg.Wait()
|
||||||
|
for _, p := range paths {
|
||||||
|
if p != "" {
|
||||||
|
defer os.Remove(p)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for i, e := range errs {
|
||||||
|
if e != "" {
|
||||||
|
t.Fatalf("fichier %d : %s", i, e)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in new issue
Block a user