From b60f441df877a649fd9e38c27820f34337ace321 Mon Sep 17 00:00:00 2001 From: Brandon Philips Date: Fri, 28 Nov 2014 09:51:31 -0500 Subject: [PATCH 1/7] cas: minor refactor, split download into download & import --- cas/remote.go | 32 ++++++++++++++++++-------------- 1 file changed, 18 insertions(+), 14 deletions(-) diff --git a/cas/remote.go b/cas/remote.go index cc0ddd6..3d0822b 100644 --- a/cas/remote.go +++ b/cas/remote.go @@ -73,24 +73,12 @@ func decompress(rs io.Reader, typ aci.FileType) (io.Reader, error) { return dr, nil } -// TODO: add locking -func (r Remote) Download(ds Store) (*Remote, error) { +func (r Remote) Import(ds Store, orig io.Reader) (*Remote, error) { var b bytes.Buffer - res, err := http.Get(r.Name) - if err != nil { - return nil, err - } - defer res.Body.Close() - - // TODO(jonboulle): handle http more robustly (redirects?) - if res.StatusCode != http.StatusOK { - return nil, fmt.Errorf("bad HTTP status code: %d", res.StatusCode) - } - // TODO(philips): use go routines to parallelize this pipeline and make // the file type detection happen without a second stream - _, err = io.Copy(&b, res.Body) + _, err := io.Copy(&b, orig) if err != nil { return nil, err } @@ -149,3 +137,19 @@ func (r Remote) Download(ds Store) (*Remote, error) { return &r, nil } + +// TODO: add locking +func (r Remote) Download(ds Store) (*Remote, error) { + res, err := http.Get(r.Name) + if err != nil { + return nil, err + } + defer res.Body.Close() + + // TODO(jonboulle): handle http more robustly (redirects?) + if res.StatusCode != http.StatusOK { + return nil, fmt.Errorf("bad HTTP status code: %d", res.StatusCode) + } + + return r.Import(ds, res.Body) +} From 11b7bf39106db4f587f96b440f18743ca7f35d5a Mon Sep 17 00:00:00 2001 From: Brandon Philips Date: Fri, 28 Nov 2014 10:02:03 -0500 Subject: [PATCH 2/7] cas: move uncompress helper into utils --- cas/remote.go | 27 --------------------------- cas/utils.go | 28 ++++++++++++++++++++++++++++ 2 files changed, 28 insertions(+), 27 deletions(-) diff --git a/cas/remote.go b/cas/remote.go index 3d0822b..3b14dd4 100644 --- a/cas/remote.go +++ b/cas/remote.go @@ -2,14 +2,11 @@ package cas import ( "bytes" - "compress/bzip2" - "compress/gzip" "crypto/sha256" "encoding/json" "fmt" "io" "net/http" - "os" "github.com/coreos-inc/rkt/app-container/aci" ) @@ -49,30 +46,6 @@ func (r Remote) Type() int64 { return remoteType } -func decompress(rs io.Reader, typ aci.FileType) (io.Reader, error) { - var ( - dr io.Reader - err error - ) - switch typ { - case aci.TypeGzip: - dr, err = gzip.NewReader(rs) - if err != nil { - return nil, err - } - case aci.TypeBzip2: - dr = bzip2.NewReader(rs) - case aci.TypeXz: - dr = aci.XzReader(rs) - case aci.TypeUnknown: - fmt.Fprintf(os.Stderr, "error: unknown image filetype\n") - default: - // should never happen - panic("no type returned from DetectFileType?") - } - return dr, nil -} - func (r Remote) Import(ds Store, orig io.Reader) (*Remote, error) { var b bytes.Buffer diff --git a/cas/utils.go b/cas/utils.go index 5feba31..fbcabbf 100644 --- a/cas/utils.go +++ b/cas/utils.go @@ -1,11 +1,16 @@ package cas import ( + "compress/bzip2" + "compress/gzip" "crypto/sha256" + "errors" "fmt" "io" "net/url" "strings" + + "github.com/coreos-inc/rkt/app-container/aci" ) // copy the default of git which is a two byte prefix. We will likely want to @@ -29,3 +34,26 @@ func parseAlways(s string) *url.URL { u, _ := url.Parse(s) return u } + +func decompress(rs io.Reader, typ aci.FileType) (io.Reader, error) { + var ( + dr io.Reader + err error + ) + switch typ { + case aci.TypeGzip: + dr, err = gzip.NewReader(rs) + if err != nil { + return nil, err + } + case aci.TypeBzip2: + dr = bzip2.NewReader(rs) + case aci.TypeXz: + dr = aci.XzReader(rs) + case aci.TypeUnknown: + return nil, errors.New("error: unknown image filetype") + default: + return nil, errors.New("no type returned from DetectFileType?") + } + return dr, nil +} From 0e2a394b57b9e94f5aba3bb6b49ca430264c42ed Mon Sep 17 00:00:00 2001 From: Brandon Philips Date: Fri, 28 Nov 2014 10:05:08 -0500 Subject: [PATCH 3/7] cas: rename the download type to tmp It is currently used as the in-process download store but it is really for any temporary stuff and whould be GC'd during the GC cycle. --- cas/cas.go | 4 ++-- cas/remote.go | 10 +++++----- 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/cas/cas.go b/cas/cas.go index c3199de..d6c2400 100644 --- a/cas/cas.go +++ b/cas/cas.go @@ -11,13 +11,13 @@ import ( const ( remoteType int64 = iota objectType - downloadType + tmpType ) var otmap = [...]string{ "remote", "object", - "download", + "tmp", } type Blob interface { diff --git a/cas/remote.go b/cas/remote.go index 3b14dd4..eeaf5fd 100644 --- a/cas/remote.go +++ b/cas/remote.go @@ -55,13 +55,13 @@ func (r Remote) Import(ds Store, orig io.Reader) (*Remote, error) { if err != nil { return nil, err } - err = ds.stores[downloadType].WriteStream(r.Hash(), &b, true) + err = ds.stores[tmpType].WriteStream(r.Hash(), &b, true) if err != nil { return nil, err } // Detect the filetype - rs, err := ds.stores[downloadType].ReadStream(r.Hash(), false) + rs, err := ds.stores[tmpType].ReadStream(r.Hash(), false) if err != nil { return nil, err } @@ -70,7 +70,7 @@ func (r Remote) Import(ds Store, orig io.Reader) (*Remote, error) { if err != nil { return nil, err } - rs, err = ds.stores[downloadType].ReadStream(r.Hash(), false) + rs, err = ds.stores[tmpType].ReadStream(r.Hash(), false) if err != nil { return nil, err } @@ -88,7 +88,7 @@ func (r Remote) Import(ds Store, orig io.Reader) (*Remote, error) { } // Store the decompressed tar - rs, err = ds.stores[downloadType].ReadStream(r.Hash(), false) + rs, err = ds.stores[tmpType].ReadStream(r.Hash(), false) if err != nil { return nil, err } @@ -104,7 +104,7 @@ func (r Remote) Import(ds Store, orig io.Reader) (*Remote, error) { return nil, err } - ds.stores[downloadType].Erase(r.Hash()) + ds.stores[tmpType].Erase(r.Hash()) r.File = key ds.stores[remoteType].Write(r.Hash(), r.Marshal()) From 77d772e1a9ec8eea200a4c6801d07be7ef717b35 Mon Sep 17 00:00:00 2001 From: Brandon Philips Date: Fri, 28 Nov 2014 10:24:49 -0500 Subject: [PATCH 4/7] cas: refactor naming of internals The indexes were called blobs and the blobs were called objects. Clean that mess up. --- cas/cas.go | 64 ++++++++++++++++++++++++++++----------------------- cas/remote.go | 6 ++--- cmd/fetch.go | 6 ++--- stage0/run.go | 2 +- 4 files changed, 42 insertions(+), 36 deletions(-) diff --git a/cas/cas.go b/cas/cas.go index d6c2400..708085b 100644 --- a/cas/cas.go +++ b/cas/cas.go @@ -8,25 +8,20 @@ import ( "github.com/coreos-inc/rkt/Godeps/_workspace/src/github.com/peterbourgon/diskv" ) +// TODO(philips): use a database for the secondary indexes like remoteType and +// appType. This is OK for now though. const ( - remoteType int64 = iota - objectType + blobType int64 = iota + remoteType tmpType ) var otmap = [...]string{ + "blob", "remote", - "object", "tmp", } -type Blob interface { - Hash() string - Marshal() []byte - Unmarshal([]byte) - Type() int64 -} - type Store struct { stores []*diskv.Diskv } @@ -46,6 +41,36 @@ func NewStore(base string) *Store { return ds } +func (ds Store) ReadStream(key string) (io.ReadCloser, error) { + return ds.stores[blobType].ReadStream(key, false) +} + +func (ds Store) WriteStream(key string, r io.Reader) error { + return ds.stores[blobType].WriteStream(key, r, true) +} + +type Index interface { + Hash() string + Marshal() []byte + Unmarshal([]byte) + Type() int64 +} + +func (ds Store) WriteIndex(i Index) { + ds.stores[i.Type()].Write(i.Hash(), i.Marshal()) +} + +func (ds Store) ReadIndex(i Index) error { + buf, err := ds.stores[i.Type()].Read(i.Hash()) + if err != nil { + return err + } + + i.Unmarshal(buf) + + return nil +} + func (ds Store) Dump(hex bool) { for _, s := range ds.stores { var keyCount int @@ -67,22 +92,3 @@ func (ds Store) Dump(hex bool) { fmt.Printf("%d total keys\n", keyCount) } } - -func (ds Store) Store(b Blob) { - ds.stores[b.Type()].Write(b.Hash(), b.Marshal()) -} - -func (ds Store) ObjectStream(file string) (io.ReadCloser, error) { - return ds.stores[objectType].ReadStream(file, false) -} - -func (ds Store) Get(b Blob) error { - buf, err := ds.stores[b.Type()].Read(b.Hash()) - if err != nil { - return err - } - - b.Unmarshal(buf) - - return nil -} diff --git a/cas/remote.go b/cas/remote.go index eeaf5fd..3d2f2b6 100644 --- a/cas/remote.go +++ b/cas/remote.go @@ -23,7 +23,7 @@ type Remote struct { Name string Mirrors []string ETag string - File string + Blob string } func (r Remote) Marshal() []byte { @@ -99,13 +99,13 @@ func (r Remote) Import(ds Store, orig io.Reader) (*Remote, error) { } key := fmt.Sprintf("sha256-%x", hash.Sum(nil)) - err = ds.stores[objectType].WriteStream(key, dr, true) + err = ds.stores[blobType].WriteStream(key, dr, true) if err != nil { return nil, err } ds.stores[tmpType].Erase(r.Hash()) - r.File = key + r.Blob = key ds.stores[remoteType].Write(r.Hash(), r.Marshal()) return &r, nil diff --git a/cmd/fetch.go b/cmd/fetch.go index eafd21c..a9c9e6a 100644 --- a/cmd/fetch.go +++ b/cmd/fetch.go @@ -23,14 +23,14 @@ var ( func fetchURL(img string, ds *cas.Store) (string, error) { rem := cas.NewRemote(img, []string{}) - err := ds.Get(rem) - if err != nil && rem.File == "" { + err := ds.ReadIndex(rem) + if err != nil && rem.Blob == "" { rem, err = rem.Download(*ds) if err != nil { return "", fmt.Errorf("downloading: %v\n", err) } } - return rem.File, nil + return rem.Blob, nil } func runFetch(args []string) (exit int) { diff --git a/stage0/run.go b/stage0/run.go index 57c2273..a41ec4e 100644 --- a/stage0/run.go +++ b/stage0/run.go @@ -264,7 +264,7 @@ func unpackBuiltinRootfs(dir string) error { func setupImage(cfg Config, img string, h types.Hash, dir string) (*schema.AppManifest, error) { log.Println("Loading image", img) - rs, err := cfg.Store.ObjectStream(img) + rs, err := cfg.Store.ReadStream(img) if err != nil { return nil, err } From 10011c6c236c4910ed9ad4a3a9b352c3e74caa1e Mon Sep 17 00:00:00 2001 From: Brandon Philips Date: Fri, 28 Nov 2014 11:06:51 -0500 Subject: [PATCH 5/7] cas: add a WriteACI method Pull the "import" out of the remote and add a WriteACI to the cas. --- cas/cas.go | 68 +++++++++++++++++++++++++++++++++++++++++++ cas/remote.go | 80 ++++++--------------------------------------------- 2 files changed, 77 insertions(+), 71 deletions(-) diff --git a/cas/cas.go b/cas/cas.go index 708085b..749c3a6 100644 --- a/cas/cas.go +++ b/cas/cas.go @@ -1,10 +1,13 @@ package cas import ( + "bytes" + "crypto/sha256" "fmt" "io" "path/filepath" + "github.com/coreos-inc/rkt/app-container/aci" "github.com/coreos-inc/rkt/Godeps/_workspace/src/github.com/peterbourgon/diskv" ) @@ -49,6 +52,71 @@ func (ds Store) WriteStream(key string, r io.Reader) error { return ds.stores[blobType].WriteStream(key, r, true) } +func (ds Store) WriteACI(tmpKey string, orig io.Reader) (string, error) { + var b bytes.Buffer + + // TODO(philips): use go routines to parallelize this pipeline and make + // the file type detection happen without a second stream + _, err := io.Copy(&b, orig) + if err != nil { + return "", err + } + err = ds.stores[tmpType].WriteStream(tmpKey, &b, true) + if err != nil { + return "", err + } + + // Detect the filetype + rs, err := ds.stores[tmpType].ReadStream(tmpKey, false) + if err != nil { + return "", err + } + defer rs.Close() + typ, err := aci.DetectFileType(rs) + if err != nil { + return "", err + } + rs, err = ds.stores[tmpType].ReadStream(tmpKey, false) + if err != nil { + return "", err + } + defer rs.Close() + + // Generate the hash of the decompressed tar + dr, err := decompress(rs, typ) + if err != nil { + return "", err + } + hash := sha256.New() + _, err = io.Copy(hash, dr) + if err != nil { + return "", err + } + + // Store the decompressed tar + rs, err = ds.stores[tmpType].ReadStream(tmpKey, false) + if err != nil { + return "", err + } + defer rs.Close() + dr, err = decompress(rs, typ) + if err != nil { + return "", err + } + + key := fmt.Sprintf("sha256-%x", hash.Sum(nil)) + err = ds.stores[blobType].WriteStream(key, dr, true) + if err != nil { + return "", err + } + + ds.stores[tmpType].Erase(tmpKey) + + return key, nil +} + + + type Index interface { Hash() string Marshal() []byte diff --git a/cas/remote.go b/cas/remote.go index 3d2f2b6..25b60af 100644 --- a/cas/remote.go +++ b/cas/remote.go @@ -1,14 +1,9 @@ package cas import ( - "bytes" - "crypto/sha256" "encoding/json" "fmt" - "io" "net/http" - - "github.com/coreos-inc/rkt/app-container/aci" ) func NewRemote(name string, mirrors []string) *Remote { @@ -46,71 +41,6 @@ func (r Remote) Type() int64 { return remoteType } -func (r Remote) Import(ds Store, orig io.Reader) (*Remote, error) { - var b bytes.Buffer - - // TODO(philips): use go routines to parallelize this pipeline and make - // the file type detection happen without a second stream - _, err := io.Copy(&b, orig) - if err != nil { - return nil, err - } - err = ds.stores[tmpType].WriteStream(r.Hash(), &b, true) - if err != nil { - return nil, err - } - - // Detect the filetype - rs, err := ds.stores[tmpType].ReadStream(r.Hash(), false) - if err != nil { - return nil, err - } - defer rs.Close() - typ, err := aci.DetectFileType(rs) - if err != nil { - return nil, err - } - rs, err = ds.stores[tmpType].ReadStream(r.Hash(), false) - if err != nil { - return nil, err - } - defer rs.Close() - - // Generate the hash of the decompressed tar - dr, err := decompress(rs, typ) - if err != nil { - return nil, err - } - hash := sha256.New() - _, err = io.Copy(hash, dr) - if err != nil { - return nil, err - } - - // Store the decompressed tar - rs, err = ds.stores[tmpType].ReadStream(r.Hash(), false) - if err != nil { - return nil, err - } - defer rs.Close() - dr, err = decompress(rs, typ) - if err != nil { - return nil, err - } - - key := fmt.Sprintf("sha256-%x", hash.Sum(nil)) - err = ds.stores[blobType].WriteStream(key, dr, true) - if err != nil { - return nil, err - } - - ds.stores[tmpType].Erase(r.Hash()) - r.Blob = key - ds.stores[remoteType].Write(r.Hash(), r.Marshal()) - - return &r, nil -} - // TODO: add locking func (r Remote) Download(ds Store) (*Remote, error) { res, err := http.Get(r.Name) @@ -124,5 +54,13 @@ func (r Remote) Download(ds Store) (*Remote, error) { return nil, fmt.Errorf("bad HTTP status code: %d", res.StatusCode) } - return r.Import(ds, res.Body) + key, err := ds.WriteACI(r.Hash(), res.Body) + if err != nil { + return nil, err + } + + r.Blob = key + ds.WriteIndex(&r) + + return &r, nil } From 44912ffdb350f9a1176822e73020fbd113077e2a Mon Sep 17 00:00:00 2001 From: Brandon Philips Date: Fri, 28 Nov 2014 11:48:22 -0500 Subject: [PATCH 6/7] app-container: introduce NewHashSHA256 Simple helper to generate a Hash type from a byte slice. --- app-container/schema/types/hash.go | 8 ++++++++ cas/remote.go | 4 +++- cas/utils.go | 8 -------- cmd/run.go | 15 +++++++++++++++ 4 files changed, 26 insertions(+), 9 deletions(-) diff --git a/app-container/schema/types/hash.go b/app-container/schema/types/hash.go index 9a7caa7..9720645 100644 --- a/app-container/schema/types/hash.go +++ b/app-container/schema/types/hash.go @@ -1,6 +1,7 @@ package types import ( + "crypto/sha256" "encoding/json" "errors" "fmt" @@ -69,3 +70,10 @@ func (h Hash) MarshalJSON() ([]byte, error) { } return json.Marshal(h.String()) } + +func NewHashSHA256(b []byte) *Hash { + h := sha256.New() + h.Write(b) + nh, _ := NewHash(fmt.Sprintf("sha256-%x", h.Sum(nil))) + return nh +} diff --git a/cas/remote.go b/cas/remote.go index 25b60af..1b53320 100644 --- a/cas/remote.go +++ b/cas/remote.go @@ -4,6 +4,8 @@ import ( "encoding/json" "fmt" "net/http" + + "github.com/coreos-inc/rkt/app-container/schema/types" ) func NewRemote(name string, mirrors []string) *Remote { @@ -34,7 +36,7 @@ func (r *Remote) Unmarshal(data []byte) { } func (r Remote) Hash() string { - return sha256sum(r.Name) + return types.NewHashSHA256([]byte(r.Name)).String() } func (r Remote) Type() int64 { diff --git a/cas/utils.go b/cas/utils.go index fbcabbf..8265cf7 100644 --- a/cas/utils.go +++ b/cas/utils.go @@ -3,9 +3,7 @@ package cas import ( "compress/bzip2" "compress/gzip" - "crypto/sha256" "errors" - "fmt" "io" "net/url" "strings" @@ -24,12 +22,6 @@ func blockTransform(s string) []string { return pathSlice } -func sha256sum(s string) string { - h := sha256.New() - io.WriteString(h, s) - return fmt.Sprintf("sha256-%x", h.Sum(nil)) -} - func parseAlways(s string) *url.URL { u, _ := url.Parse(s) return u diff --git a/cmd/run.go b/cmd/run.go index 4a3d322..1555b43 100644 --- a/cmd/run.go +++ b/cmd/run.go @@ -43,6 +43,21 @@ func findImages(args []string, ds *cas.Store) (out []string, err error) { if err == nil { continue } + + // import the local file if it exists + file, err := os.Open(img) + if err == nil { + hash := types.NewHashSHA256([]byte(img)).String() + key, err := ds.WriteACI(hash, file) + file.Close() + if err != nil { + return nil, fmt.Errorf("%s: %v", img, err) + } + out[i] = key + continue + } + + // download if it is a URL u, err := url.Parse(img) if err != nil { return nil, fmt.Errorf("%s: not a valid URL or hash", img) From 883c12ba3457c848b1265fe9facbaefa237ed63c Mon Sep 17 00:00:00 2001 From: Brandon Philips Date: Fri, 28 Nov 2014 12:46:30 -0500 Subject: [PATCH 7/7] cmd: run: introduce a helpful description about IMAGE --- cmd/run.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/cmd/run.go b/cmd/run.go index 1555b43..2b468bd 100644 --- a/cmd/run.go +++ b/cmd/run.go @@ -23,6 +23,8 @@ var ( Name: "run", Summary: "Run image(s) in an application container in rocket", Usage: "[--volume LABEL:SOURCE] IMAGE...", + Description: `IMAGE should be a string referencing an image; either a hash, local file on disk, or URL. +They will be checked in that order and the first match will be used.`, Run: runRun, } )