Keep a facts snapshot for merge checks, and say when it goes stale (hq to-be 45 Phase 5, S14)
Every check the mesh had was right about the world it was given and none was given the mesh's: a real machine's name made an identity too long (263), the node-engine refused what the catalogue check passed (236). The controller now composes what a check needs - every machine under a pseudonym of its name's length, its roles, system, builds, capabilities, assignments, pins, settings and how its declaration composes; every seat, module and source; the bus, store and node-engine versions it runs - with no secret, no address and no name, and keeps it in the artifact store as facts:latest when it moved, or daily. The replaced snapshot's manifest is let go of, so the nightly collector takes it. S14 raises facts-stale past two days.
This commit is contained in:
@@ -0,0 +1,179 @@
|
||||
package artifacts
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// A document the mesh keeps under a name, read by whoever asks for that name (novox/hq to-be 45 §9).
|
||||
//
|
||||
// **The one thing the mesh names by tag.** Everything a machine runs is pinned by digest (ADR 0189), and
|
||||
// holders are put untagged so the collector's own rule keeps them. The facts snapshot is the opposite
|
||||
// case: what a reader wants is *the newest*, never a particular one, and a tag is the store's own word
|
||||
// for that. So it is put as the smallest OCI manifest naming one layer, under a tag; putting the next
|
||||
// moves the tag, and **the manifest it replaced is deleted by its digest**: the store's nightly collector
|
||||
// keeps every manifest, tagged or not, and marks what each names — so a replaced snapshot left in place
|
||||
// would be kept for ever, one more every day. Deleted, its layer is named by nothing and the next
|
||||
// collection takes it: the store keeps the newest, and nothing grows.
|
||||
|
||||
// MaxTagged is the largest document put or read this way: a snapshot of a large mesh is a few
|
||||
// megabytes, and a reader is never made to swallow an answer of any size.
|
||||
const MaxTagged = 64 << 20
|
||||
|
||||
// ErrNoTag is the answer when the store holds nothing under the tag: never put, or collected.
|
||||
var ErrNoTag = errors.New("the artifact store holds nothing under that name")
|
||||
|
||||
// PutTagged puts body as the only layer of a manifest in repository and points tag at it. Answers the
|
||||
// layer's digest — what a reader quotes as "the snapshot I read".
|
||||
func (s Store) PutTagged(ctx context.Context, repository, tag, mediaType string, body []byte) (string, error) {
|
||||
if s.Address == "" {
|
||||
return "", fmt.Errorf("this mesh has no artifact store on its network to keep %s:%s in", repository, tag)
|
||||
}
|
||||
if len(body) > MaxTagged {
|
||||
return "", fmt.Errorf("%s:%s is %d bytes, over the %d the store is given", repository, tag, len(body), MaxTagged)
|
||||
}
|
||||
sum := sha256.Sum256(body)
|
||||
digest := "sha256:" + hex.EncodeToString(sum[:])
|
||||
// What the tag names now, so it can be let go of once the new one stands. Asked before the put:
|
||||
// after it, the tag names the new one.
|
||||
replaced, err := s.manifestDigest(ctx, repository, tag)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := s.putBlob(ctx, repository, digest, body); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := s.putBlob(ctx, repository, emptyDigest, emptyConfig); err != nil {
|
||||
return "", err
|
||||
}
|
||||
manifest, err := json.Marshal(holderManifest{
|
||||
SchemaVersion: 2,
|
||||
MediaType: mediaManifest,
|
||||
Config: descriptor{MediaType: mediaEmpty, Digest: emptyDigest, Size: int64(len(emptyConfig))},
|
||||
Layers: []descriptor{{MediaType: mediaType, Digest: digest, Size: int64(len(body))}},
|
||||
})
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
request, err := http.NewRequestWithContext(ctx, http.MethodPut, s.url(repository, "manifests", tag),
|
||||
bytes.NewReader(manifest))
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
request.Header.Set("Content-Type", mediaManifest)
|
||||
response, err := s.client().Do(request)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("cannot reach the artifact store at %s: %w", s.Address, err)
|
||||
}
|
||||
defer response.Body.Close()
|
||||
if response.StatusCode != http.StatusCreated {
|
||||
said, _ := io.ReadAll(io.LimitReader(response.Body, 4096))
|
||||
return "", fmt.Errorf("the artifact store refused %s:%s: %s %s", repository, tag, response.Status,
|
||||
strings.TrimSpace(string(said)))
|
||||
}
|
||||
mSum := sha256.Sum256(manifest)
|
||||
if put := "sha256:" + hex.EncodeToString(mSum[:]); replaced != "" && replaced != put {
|
||||
// The new one stands; the old one is let go of. A refusal here leaves one more snapshot in the
|
||||
// store, which is said and is not a failure of the put: the tag already names the new one.
|
||||
if err := s.remove(ctx, s.url(repository, "manifests", replaced), repository+"/manifests/"+replaced); err != nil &&
|
||||
err != Gone {
|
||||
return digest, fmt.Errorf("%s:%s now names the new document, and the one it replaced could not be "+
|
||||
"let go of: %w", repository, tag, err)
|
||||
}
|
||||
}
|
||||
return digest, nil
|
||||
}
|
||||
|
||||
// manifestDigest is the digest of the manifest tag names in repository; empty when it names none.
|
||||
func (s Store) manifestDigest(ctx context.Context, repository, tag string) (string, error) {
|
||||
request, err := http.NewRequestWithContext(ctx, http.MethodHead, s.url(repository, "manifests", tag), nil)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
for _, media := range manifestAccept {
|
||||
request.Header.Add("Accept", media)
|
||||
}
|
||||
response, err := s.client().Do(request)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("cannot reach the artifact store at %s: %w", s.Address, err)
|
||||
}
|
||||
defer response.Body.Close()
|
||||
switch response.StatusCode {
|
||||
case http.StatusOK:
|
||||
return response.Header.Get("Docker-Content-Digest"), nil
|
||||
case http.StatusNotFound:
|
||||
return "", nil
|
||||
default:
|
||||
return "", fmt.Errorf("the artifact store answered %s for %s:%s", response.Status, repository, tag)
|
||||
}
|
||||
}
|
||||
|
||||
// GetTagged reads the one layer of the manifest tag names in repository, and its digest. ErrNoTag when
|
||||
// the store holds nothing under it.
|
||||
func (s Store) GetTagged(ctx context.Context, repository, tag string) ([]byte, string, error) {
|
||||
if s.Address == "" {
|
||||
return nil, "", fmt.Errorf("this mesh has no artifact store on its network to read %s:%s from", repository, tag)
|
||||
}
|
||||
request, err := http.NewRequestWithContext(ctx, http.MethodGet, s.url(repository, "manifests", tag), nil)
|
||||
if err != nil {
|
||||
return nil, "", err
|
||||
}
|
||||
for _, media := range manifestAccept {
|
||||
request.Header.Add("Accept", media)
|
||||
}
|
||||
response, err := s.client().Do(request)
|
||||
if err != nil {
|
||||
return nil, "", fmt.Errorf("cannot reach the artifact store at %s: %w", s.Address, err)
|
||||
}
|
||||
defer response.Body.Close()
|
||||
switch response.StatusCode {
|
||||
case http.StatusOK:
|
||||
case http.StatusNotFound:
|
||||
return nil, "", fmt.Errorf("%w: %s:%s", ErrNoTag, repository, tag)
|
||||
default:
|
||||
return nil, "", fmt.Errorf("the artifact store answered %s for %s:%s", response.Status, repository, tag)
|
||||
}
|
||||
var m holderManifest
|
||||
if err := json.NewDecoder(io.LimitReader(response.Body, 1<<20)).Decode(&m); err != nil {
|
||||
return nil, "", fmt.Errorf("%s:%s is not a manifest the mesh wrote: %w", repository, tag, err)
|
||||
}
|
||||
if len(m.Layers) != 1 {
|
||||
return nil, "", fmt.Errorf("%s:%s names %d layers; the mesh writes one", repository, tag, len(m.Layers))
|
||||
}
|
||||
layer := m.Layers[0]
|
||||
if layer.Size > MaxTagged {
|
||||
return nil, "", fmt.Errorf("%s:%s is %d bytes, over the %d a reader takes", repository, tag, layer.Size, MaxTagged)
|
||||
}
|
||||
blob, err := http.NewRequestWithContext(ctx, http.MethodGet, s.url(repository, "blobs", layer.Digest), nil)
|
||||
if err != nil {
|
||||
return nil, "", err
|
||||
}
|
||||
got, err := s.client().Do(blob)
|
||||
if err != nil {
|
||||
return nil, "", fmt.Errorf("cannot reach the artifact store at %s: %w", s.Address, err)
|
||||
}
|
||||
defer got.Body.Close()
|
||||
if got.StatusCode != http.StatusOK {
|
||||
return nil, "", fmt.Errorf("the artifact store names %s for %s:%s and answered %s for it",
|
||||
layer.Digest, repository, tag, got.Status)
|
||||
}
|
||||
body, err := io.ReadAll(io.LimitReader(got.Body, MaxTagged+1))
|
||||
if err != nil {
|
||||
return nil, "", err
|
||||
}
|
||||
// **Read back against its own digest**: a reader is told which snapshot it read, and a truncated or
|
||||
// substituted body must not pass as that one.
|
||||
sum := sha256.Sum256(body)
|
||||
if "sha256:"+hex.EncodeToString(sum[:]) != layer.Digest {
|
||||
return nil, "", fmt.Errorf("%s:%s read back as something other than %s", repository, tag, layer.Digest)
|
||||
}
|
||||
return body, layer.Digest, nil
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
package artifacts
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// The facts snapshot is put under a tag and read back by it (novox/hq to-be 45 §9).
|
||||
//
|
||||
// Against the very registry the mesh's store runs, because what is asserted is the registry's answer:
|
||||
// that a manifest put by tag is read by tag, that the layer comes back whole and checked, and that the
|
||||
// snapshot a newer one replaced is the collector's to take while the newest is kept. Raised as the
|
||||
// collector test above says, with MESH_TEST_REGISTRY and MESH_TEST_REGISTRY_CONTAINER.
|
||||
func TestLiveADocumentPutUnderATagIsReadBackAndOnlyTheNewestIsKept(t *testing.T) {
|
||||
address := os.Getenv("MESH_TEST_REGISTRY")
|
||||
container := os.Getenv("MESH_TEST_REGISTRY_CONTAINER")
|
||||
if address == "" || container == "" {
|
||||
t.Skip("no MESH_TEST_REGISTRY / MESH_TEST_REGISTRY_CONTAINER; see the collector test's comment for the registry to raise")
|
||||
}
|
||||
ctx := context.Background()
|
||||
store := Store{Address: address}
|
||||
repository := fmt.Sprintf("facts-live-%d", time.Now().UnixNano())
|
||||
|
||||
if _, _, err := store.GetTagged(ctx, repository, "latest"); !errors.Is(err, ErrNoTag) {
|
||||
t.Fatalf("nothing was put and the read said %v, not that nothing is there", err)
|
||||
}
|
||||
first := []byte(`{"facts":1,"taken":"first"}`)
|
||||
firstDigest, err := store.PutTagged(ctx, repository, "latest", "application/vnd.novox.mesh.facts.v1+json", first)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
second := []byte(`{"facts":1,"taken":"second"}`)
|
||||
secondDigest, err := store.PutTagged(ctx, repository, "latest", "application/vnd.novox.mesh.facts.v1+json", second)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got, digest, err := store.GetTagged(ctx, repository, "latest")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if string(got) != string(second) || digest != secondDigest {
|
||||
t.Fatalf("read back %q (%s), not the newest put %q (%s)", got, digest, second, secondDigest)
|
||||
}
|
||||
|
||||
// The collector, as the store's nightly step runs it — with no `--delete-untagged`: the newest kept,
|
||||
// the replaced one taken because its manifest was let go of.
|
||||
out, err := exec.Command("docker", "exec", container, "registry", "garbage-collect",
|
||||
"/etc/docker/registry/config.yml").CombinedOutput()
|
||||
if err != nil {
|
||||
t.Fatalf("the collector did not run: %v\n%s", err, out)
|
||||
}
|
||||
if got, _, err := store.GetTagged(ctx, repository, "latest"); err != nil || string(got) != string(second) {
|
||||
t.Fatalf("after collection the newest reads %q, %v", got, err)
|
||||
}
|
||||
// On the store's disk, not as the running server answers: it caches blob descriptors in memory.
|
||||
onDisk := func(digest string) bool {
|
||||
hex := strings.TrimPrefix(digest, "sha256:")
|
||||
path := "/var/lib/registry/docker/registry/v2/blobs/sha256/" + hex[:2] + "/" + hex + "/data"
|
||||
return exec.Command("docker", "exec", container, "test", "-f", path).Run() == nil
|
||||
}
|
||||
if onDisk(firstDigest) {
|
||||
t.Errorf("the replaced snapshot %s is still in the store after collection; nothing would ever take it", firstDigest)
|
||||
}
|
||||
if !onDisk(secondDigest) {
|
||||
t.Errorf("the collector took the newest snapshot %s", secondDigest)
|
||||
}
|
||||
if !strings.HasPrefix(secondDigest, "sha256:") {
|
||||
t.Errorf("the digest said %q", secondDigest)
|
||||
}
|
||||
}
|
||||
|
||||
// A store with no address is said, not dialled.
|
||||
func TestADocumentWithNowhereToGoIsRefusedByName(t *testing.T) {
|
||||
if _, err := (Store{}).PutTagged(context.Background(), "facts", "latest", "x", []byte("{}")); err == nil ||
|
||||
!strings.Contains(err.Error(), "no artifact store") {
|
||||
t.Errorf("put with no store said %v", err)
|
||||
}
|
||||
if _, _, err := (Store{}).GetTagged(context.Background(), "facts", "latest"); err == nil ||
|
||||
!strings.Contains(err.Error(), "no artifact store") {
|
||||
t.Errorf("read with no store said %v", err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user