diff --git a/script/qemu_cache/classify.go b/script/qemu_cache/classify.go index 5132a93..5592338 100644 --- a/script/qemu_cache/classify.go +++ b/script/qemu_cache/classify.go @@ -98,6 +98,20 @@ var volatileSuffixes = []string{ var parEmpreinte = regexp.MustCompile( `/by-hash/(MD5Sum|SHA1|SHA256|SHA512)/[0-9a-fA-F]{32,128}$`) +// parEmpreinteDeDepot reconnaît une métadonnée de dépôt RPM nommée par +// l'empreinte de son contenu : « …/repodata/-primary.xml.zck ». +// Seul « repomd.xml » y est volatile : il désigne la version COURANTE de +// chaque métadonnée par ce nom, si bien qu'un contenu nouveau porte un nom +// nouveau. Le fichier nommé est aussi figé qu'un paquet. +// +// Volatile, il est repris en entier à chaque installation — un « primary » +// pèse des dizaines de mégaoctets par dépôt. Et dnf télécharge un « .zck » par +// plages, qu'aucun volatile ne garde : amont coupé, la VM suivante n'avait +// rien. Figé, il sort du disque, et une plage y déclenche la prise du fichier +// entier (voir completer). +var parEmpreinteDeDepot = regexp.MustCompile( + `/repodata/[0-9a-fA-F]{32,128}-[^/]+$`) + // dernierePublication reconnaît « ///releases/latest/ // download/ » : un POINTEUR vers la dernière version publiée, dont // la cible change à chaque publication. @@ -165,7 +179,7 @@ func Classify(u *url.URL) Class { if dernierePublication.MatchString(u.Path) { return ClassVolatile } - if parEmpreinte.MatchString(u.Path) { + if parEmpreinte.MatchString(u.Path) || parEmpreinteDeDepot.MatchString(u.Path) { return ClassImmutable } for _, n := range volatileNames { @@ -207,7 +221,7 @@ func PortableParChemin(u *url.URL) bool { // POSITIVE parce qu'aucune des tables suivantes ne le reconnaîtrait : une // somme hexadécimale n'a pas d'extension, et retirer la seule exclusion ne // suffirait donc pas à le rendre portable. - if parEmpreinte.MatchString(u.Path) { + if parEmpreinte.MatchString(u.Path) || parEmpreinteDeDepot.MatchString(u.Path) { return true } name := strings.ToLower(path.Base(u.Path)) diff --git a/script/qemu_cache/classify_test.go b/script/qemu_cache/classify_test.go index b1a1202..7e359c3 100644 --- a/script/qemu_cache/classify_test.go +++ b/script/qemu_cache/classify_test.go @@ -209,6 +209,36 @@ func TestLesIndexParEmpreinteSontImmuables(t *testing.T) { } } +// Une métadonnée RPM nommée par sa somme est figée, quelle que soit sa +// compression ; « repomd.xml », qui la désigne, reste volatile, et un nom sans +// somme en tête garde la règle de son suffixe. +func TestLesMetadonneesRPMParEmpreinteSontImmuables(t *testing.T) { + somme := strings.Repeat("0123456789abcdef", 4) + for _, nom := range []string{ + somme + "-primary.xml.zck", somme + "-primary.xml.gz", + somme + "-updateinfo.xml.zst", somme + "-comps-BaseOS.x86_64.xml", + } { + u, _ := url.Parse("https://miroir.example/fedora/linux/updates/42/Everything/x86_64/repodata/" + nom) + if got := Classify(u); got != ClassImmutable { + t.Errorf("%s classé « %s », attendu « immutable »", nom, got) + } + if !PortableParChemin(u) { + t.Errorf("%s n'est pas jugé portable", nom) + } + } + for _, brut := range []string{ + "https://miroir.example/fedora/repodata/repomd.xml", + "https://miroir.example/fedora/repodata/primary.xml.gz", + "https://miroir.example/fedora/repodata/pas-une-somme-primary.xml.zck", + "https://miroir.example/fedora/ailleurs/" + somme + "-primary.xml.zck", + } { + u, _ := url.Parse(brut) + if got := Classify(u); got != ClassVolatile { + t.Errorf("%s classé « %s », attendu « volatile »", brut, got) + } + } +} + // Le même index par empreinte, servi par deux miroirs. Sans clé portable, // changer de miroir vide le cache de ses index : une installation hors ligne // échoue alors sur des octets que le magasin détient pourtant, et le message diff --git a/script/qemu_cache/completion_test.go b/script/qemu_cache/completion_test.go new file mode 100644 index 0000000..b16db6e --- /dev/null +++ b/script/qemu_cache/completion_test.go @@ -0,0 +1,139 @@ +// © 2026 TechnoLibre (http://www.technolibre.ca) +// License AGPL-3.0 or later (http://www.gnu.org/licenses/agpl) + +package main + +import ( + "bytes" + "net/http" + "net/http/httptest" + "net/url" + "strings" + "sync/atomic" + "testing" + "time" +) + +// amontAPlages sert un corps fixe et honore « Range » comme un miroir : 206 +// et le seul fragment demandé. Il compte séparément les plages et les corps +// entiers envoyés. +type amontAPlages struct { + srv *httptest.Server + plages int64 + entiers int64 +} + +func nouvelAmontAPlages(t *testing.T, corps string) *amontAPlages { + t.Helper() + a := &amontAPlages{} + a.srv = httptest.NewServer(http.HandlerFunc( + func(w http.ResponseWriter, r *http.Request) { + if r.Header.Get("Range") != "" { + atomic.AddInt64(&a.plages, 1) + } else { + atomic.AddInt64(&a.entiers, 1) + } + http.ServeContent(w, r, "", time.Time{}, bytes.NewReader([]byte(corps))) + })) + t.Cleanup(a.srv.Close) + return a +} + +func (a *amontAPlages) hote() string { + return strings.TrimPrefix(a.srv.URL, "http://") +} + +func demandePlage(t *testing.T, p *Proxy, hote, chemin, plage string) *httptest.ResponseRecorder { + t.Helper() + r := httptest.NewRequest("GET", chemin, nil) + r.Host = hote + r.Header.Set("Range", plage) + w := httptest.NewRecorder() + p.serve(w, r, "http") + return w +} + +const metadonneeZck = "/fedora/updates/repodata/" + + "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef" + + "-primary.xml.zck" + +// Une plage sur une métadonnée figée absente fait prendre le fichier entier, +// une fois ; les plages suivantes sortent du disque sans rien demander à +// l'amont. +func TestUnePlageFaitGarderLeFichierEntier(t *testing.T) { + a := nouvelAmontAPlages(t, "en-tete|morceau-1|morceau-2") + p := proxyDeTest(t) + + w := demandePlage(t, p, a.hote(), metadonneeZck, "bytes=0-6") + if w.Code != http.StatusPartialContent || w.Body.String() != "en-tete" { + t.Fatalf("première plage : %d %q", w.Code, w.Body.String()) + } + p.attendreCompletions() + + w = demandePlage(t, p, a.hote(), metadonneeZck, "bytes=8-16") + if w.Code != http.StatusPartialContent || w.Body.String() != "morceau-1" { + t.Errorf("seconde plage : %d %q", w.Code, w.Body.String()) + } + if got := w.Header().Get("X-ERPLibre-Cache"); got != OutcomeHit { + t.Errorf("seconde plage servie « %s », attendu « %s »", got, OutcomeHit) + } + if n := atomic.LoadInt64(&a.entiers); n != 1 { + t.Errorf("%d corps entiers pris à l'amont, attendu 1", n) + } + if n := atomic.LoadInt64(&a.plages); n != 1 { + t.Errorf("%d plages envoyées à l'amont, attendu 1", n) + } +} + +// Amont coupé, la plage d'une VM suivante sort du disque : c'est ce qui manquait +// à une installation hors ligne. +func TestHorsLigneLaPlageSortDuDisque(t *testing.T) { + a := nouvelAmontAPlages(t, "en-tete|morceau-1|morceau-2") + p := proxyDeTest(t) + hote := a.hote() + demandePlage(t, p, hote, metadonneeZck, "bytes=0-6") + p.attendreCompletions() + a.srv.Close() + + w := demandePlage(t, p, hote, metadonneeZck, "bytes=18-26") + if w.Code != http.StatusPartialContent || w.Body.String() != "morceau-2" { + t.Errorf("hors ligne : %d %q", w.Code, w.Body.String()) + } +} + +// Un index volatile demandé par plage ne déclenche aucune prise : il n'est +// jamais servi du disque tant que l'amont répond, le garder n'épargnerait rien. +func TestUnePlageSurUnIndexNePrendRien(t *testing.T) { + a := nouvelAmontAPlages(t, "index de dépôt") + p := proxyDeTest(t) + + demandePlage(t, p, a.hote(), "/fedora/repodata/repomd.xml", "bytes=0-4") + p.attendreCompletions() + if n := atomic.LoadInt64(&a.entiers); n != 0 { + t.Errorf("%d corps entiers pris pour un index volatile, attendu 0", n) + } +} + +// Plusieurs plages simultanées sur le même fichier ne prennent le corps entier +// qu'une fois. +func TestDesPlagesSimultaneesNePrennentQuUneFois(t *testing.T) { + a := nouvelAmontAPlages(t, strings.Repeat("x", 1<<16)) + p := proxyDeTest(t) + u, _ := url.Parse("http://" + a.hote() + metadonneeZck) + key := CleDe("GET", u) + p.completions.Store(key, true) // une prise est déjà en cours + + for i := 0; i < 5; i++ { + demandePlage(t, p, a.hote(), metadonneeZck, "bytes=0-9") + } + p.attendreCompletions() + if n := atomic.LoadInt64(&a.entiers); n != 0 { + t.Errorf("%d prises lancées alors qu'une était en cours, attendu 0", n) + } + p.completions.Delete(key) + demandePlage(t, p, a.hote(), metadonneeZck, "bytes=0-9") + p.attendreCompletions() + if n := atomic.LoadInt64(&a.entiers); n != 1 { + t.Errorf("%d prises après la fin de la première, attendu 1", n) + } +} diff --git a/script/qemu_cache/main.go b/script/qemu_cache/main.go index ab495ce..c623909 100644 --- a/script/qemu_cache/main.go +++ b/script/qemu_cache/main.go @@ -25,7 +25,7 @@ import ( "time" ) -const version = "0.2.9" +const version = "0.2.10" func main() { var ( diff --git a/script/qemu_cache/proxy.go b/script/qemu_cache/proxy.go index 52cf813..fbe771a 100644 --- a/script/qemu_cache/proxy.go +++ b/script/qemu_cache/proxy.go @@ -4,8 +4,10 @@ package main import ( + "context" "encoding/json" "fmt" + "io" "log" "net" "net/http" @@ -134,6 +136,11 @@ type Proxy struct { // vise l'un d'eux sur une adresse de cette machine est une boucle. Vide, // rien n'est refusé. Ecoutes []int + + // completions tient les clés dont le corps entier est en cours de prise + // (voir completer) ; enCours les compte, pour qui doit les attendre. + completions sync.Map + enCours sync.WaitGroup } // NewProxy monte le client amont. Aucun délai GLOBAL n'est posé : une image @@ -244,6 +251,11 @@ func (p *Proxy) serve(w http.ResponseWriter, r *http.Request, scheme string) { if p.serveFromStore(w, r, u, key, class, OutcomeHit) { return } + // La plage passe à l'amont et ne se garde pas ; le fichier entier est + // pris à part, pour que la plage suivante sorte du disque. + if partial && r.Method == "GET" { + p.completer(r, u, key, class) + } } // La négociation git n'est pas relayée quand un miroir peut la servir : @@ -540,6 +552,88 @@ func varieSurAccept(h http.Header) bool { return false } +// delaiCompletion borne la prise d'un corps entier en arrière-plan. Aucun +// client n'attend cette prise : sans borne, un amont qui cesse d'envoyer au +// milieu du corps la garderait ouverte pour toujours. +const delaiCompletion = time.Hour + +// completer prend à l'amont, en arrière-plan, le corps ENTIER d'un fichier +// figé dont un client n'a demandé qu'une plage, et le garde sous sa clé. +// +// dnf télécharge ses métadonnées zchunk par plages, et pacman reprend de même +// un paquet interrompu : une plage ne se garde pas, si bien que sans cette +// prise ces fichiers repartiraient à l'amont à chaque VM — et, amont coupé, ne +// seraient pas là. Une fois gardé, le fichier sert toute plage depuis le +// disque. +// +// Une seule prise par clé à la fois, et la requête du client ne l'attend pas. +// Elle part sans plage ni condition, avec le seul agent de l'invité, et un +// amont connu muet n'est pas recomposé. Une prise manquée se retente à la +// plage suivante. +func (p *Proxy) completer(r *http.Request, u *url.URL, key string, class Class) { + if _, deja := p.completions.LoadOrStore(key, true); deja { + return + } + entete := http.Header{} + if agent := r.Header.Get("User-Agent"); agent != "" { + entete.Set("User-Agent", agent) + } + p.enCours.Add(1) + go func() { + defer p.enCours.Done() + defer p.completions.Delete(key) + p.prendreEntier(u, key, class, entete) + }() +} + +// attendreCompletions rend la main quand plus aucune prise n'est en cours. +func (p *Proxy) attendreCompletions() { + p.enCours.Wait() +} + +// prendreEntier fait la prise elle-même : un 200 entier, publié sous la clé +// par le même écrivain que le chemin du client, taille annoncée vérifiée. Tout +// autre statut, ou une clé devenue détenue entre-temps, ne publie rien. +func (p *Proxy) prendreEntier(u *url.URL, key string, class Class, entete http.Header) { + ctx, annuler := context.WithTimeout(context.Background(), delaiCompletion) + defer annuler() + req, err := http.NewRequestWithContext(ctx, "GET", u.String(), nil) + if err != nil { + return + } + req.Header = entete + resp, err := p.fetch(req, u, true) + if err != nil { + return + } + resp = p.suivreRedirections(req, u, resp) + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK || p.Store.Detient(key) { + return + } + cw, err := p.Store.NewWriter(key, Meta{ + URL: u.String(), Method: "GET", Status: resp.StatusCode, + Header: resp.Header.Clone(), Class: class.String(), + }) + if err != nil { + log.Printf(T("cache : écriture impossible pour %s : %v"), u, err) + return + } + n, err := io.Copy(cw, resp.Body) + if err != nil { + cw.Abort() + return + } + if err := cw.Commit(resp.ContentLength); err != nil { + log.Printf(T("cache : %s non gardé : %v"), u, err) + return + } + p.record(accessLine{ + Method: "GET", URL: u.String(), Class: class.String(), + Outcome: OutcomeStored, Status: resp.StatusCode, Bytes: n, Upstream: true, + }) +} + // etagGarde rend l'ETag du corps gardé sous la clé, ou "" : rien de gardé, // un statut seul, ou une réponse d'amont qui n'en portait pas. func (p *Proxy) etagGarde(key string) string { diff --git a/script/qemu_cache/proxy_test.go b/script/qemu_cache/proxy_test.go index 8de0aa8..ca2be28 100644 --- a/script/qemu_cache/proxy_test.go +++ b/script/qemu_cache/proxy_test.go @@ -161,26 +161,22 @@ func TestDefautHorsLigneNommeLeFichier(t *testing.T) { } } -// Une requête partielle ne remplit pas le cache : un fragment ne sert à rien -// à la demande suivante, et le garder comme un corps entier servirait un -// paquet tronqué. +// Une requête partielle ne garde jamais SON fragment : le garder comme un +// corps entier servirait un paquet tronqué. Ce qui entre au cache est le corps +// entier, pris à part (voir completion_test.go). func TestRequetePartielleNonGardee(t *testing.T) { - a := nouvelAmont(t, "0123456789") + a := nouvelAmontAPlages(t, "0123456789") p := proxyDeTest(t) chemin := "/x/paquet-1-1-x86_64.pkg.tar.zst" - r := httptest.NewRequest("GET", chemin, nil) - r.Host = a.hote() - r.Header.Set("Range", "bytes=0-4") - p.serve(httptest.NewRecorder(), r, "http") - - // Une demande entière ensuite doit ressortir à l'amont. - w := demande(t, p, a.hote(), chemin) - if got := w.Header().Get("X-ERPLibre-Cache"); got == OutcomeHit { - t.Error("un fragment est entré au cache et a été servi comme un corps entier") + if w := demandePlage(t, p, a.hote(), chemin, "bytes=0-4"); w.Body.String() != "01234" { + t.Fatalf("plage : %q", w.Body.String()) } - if n := a.appels(); n != 2 { - t.Errorf("l'amont a reçu %d requêtes, attendu 2", n) + p.attendreCompletions() + + w := demande(t, p, a.hote(), chemin) + if w.Body.String() != "0123456789" { + t.Errorf("demande entière servie %q : un fragment a été gardé", w.Body.String()) } }