mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-19 02:12:23 +00:00
swarm: propagate ctx to internal apis (#754)
This commit is contained in:
parent
c9b3f9e675
commit
4af94a2962
24 changed files with 170 additions and 159 deletions
|
|
@ -18,6 +18,7 @@
|
||||||
package main
|
package main
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"os"
|
"os"
|
||||||
|
|
||||||
|
|
@ -39,7 +40,7 @@ func hash(ctx *cli.Context) {
|
||||||
|
|
||||||
stat, _ := f.Stat()
|
stat, _ := f.Stat()
|
||||||
fileStore := storage.NewFileStore(storage.NewMapChunkStore(), storage.NewFileStoreParams())
|
fileStore := storage.NewFileStore(storage.NewMapChunkStore(), storage.NewFileStoreParams())
|
||||||
addr, _, err := fileStore.Store(f, stat.Size(), false)
|
addr, _, err := fileStore.Store(context.TODO(), f, stat.Size(), false)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
utils.Fatalf("%v\n", err)
|
utils.Fatalf("%v\n", err)
|
||||||
} else {
|
} else {
|
||||||
|
|
|
||||||
|
|
@ -227,28 +227,28 @@ func NewAPI(fileStore *storage.FileStore, dns Resolver, resourceHandler *mru.Han
|
||||||
}
|
}
|
||||||
|
|
||||||
// Upload to be used only in TEST
|
// Upload to be used only in TEST
|
||||||
func (a *API) Upload(uploadDir, index string, toEncrypt bool) (hash string, err error) {
|
func (a *API) Upload(ctx context.Context, uploadDir, index string, toEncrypt bool) (hash string, err error) {
|
||||||
fs := NewFileSystem(a)
|
fs := NewFileSystem(a)
|
||||||
hash, err = fs.Upload(uploadDir, index, toEncrypt)
|
hash, err = fs.Upload(uploadDir, index, toEncrypt)
|
||||||
return hash, err
|
return hash, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// Retrieve FileStore reader API
|
// Retrieve FileStore reader API
|
||||||
func (a *API) Retrieve(addr storage.Address) (reader storage.LazySectionReader, isEncrypted bool) {
|
func (a *API) Retrieve(ctx context.Context, addr storage.Address) (reader storage.LazySectionReader, isEncrypted bool) {
|
||||||
return a.fileStore.Retrieve(addr)
|
return a.fileStore.Retrieve(ctx, addr)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Store wraps the Store API call of the embedded FileStore
|
// Store wraps the Store API call of the embedded FileStore
|
||||||
func (a *API) Store(data io.Reader, size int64, toEncrypt bool) (addr storage.Address, wait func(), err error) {
|
func (a *API) Store(ctx context.Context, data io.Reader, size int64, toEncrypt bool) (addr storage.Address, wait func(), err error) {
|
||||||
log.Debug("api.store", "size", size)
|
log.Debug("api.store", "size", size)
|
||||||
return a.fileStore.Store(data, size, toEncrypt)
|
return a.fileStore.Store(ctx, data, size, toEncrypt)
|
||||||
}
|
}
|
||||||
|
|
||||||
// ErrResolve is returned when an URI cannot be resolved from ENS.
|
// ErrResolve is returned when an URI cannot be resolved from ENS.
|
||||||
type ErrResolve error
|
type ErrResolve error
|
||||||
|
|
||||||
// Resolve resolves a URI to an Address using the MultiResolver.
|
// Resolve resolves a URI to an Address using the MultiResolver.
|
||||||
func (a *API) Resolve(uri *URI) (storage.Address, error) {
|
func (a *API) Resolve(ctx context.Context, uri *URI) (storage.Address, error) {
|
||||||
apiResolveCount.Inc(1)
|
apiResolveCount.Inc(1)
|
||||||
log.Trace("resolving", "uri", uri.Addr)
|
log.Trace("resolving", "uri", uri.Addr)
|
||||||
|
|
||||||
|
|
@ -286,17 +286,17 @@ func (a *API) Resolve(uri *URI) (storage.Address, error) {
|
||||||
}
|
}
|
||||||
|
|
||||||
// Put provides singleton manifest creation on top of FileStore store
|
// Put provides singleton manifest creation on top of FileStore store
|
||||||
func (a *API) Put(content, contentType string, toEncrypt bool) (k storage.Address, wait func(), err error) {
|
func (a *API) Put(ctx context.Context, content string, contentType string, toEncrypt bool) (k storage.Address, wait func(), err error) {
|
||||||
apiPutCount.Inc(1)
|
apiPutCount.Inc(1)
|
||||||
r := strings.NewReader(content)
|
r := strings.NewReader(content)
|
||||||
key, waitContent, err := a.fileStore.Store(r, int64(len(content)), toEncrypt)
|
key, waitContent, err := a.fileStore.Store(ctx, r, int64(len(content)), toEncrypt)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiPutFail.Inc(1)
|
apiPutFail.Inc(1)
|
||||||
return nil, nil, err
|
return nil, nil, err
|
||||||
}
|
}
|
||||||
manifest := fmt.Sprintf(`{"entries":[{"hash":"%v","contentType":"%s"}]}`, key, contentType)
|
manifest := fmt.Sprintf(`{"entries":[{"hash":"%v","contentType":"%s"}]}`, key, contentType)
|
||||||
r = strings.NewReader(manifest)
|
r = strings.NewReader(manifest)
|
||||||
key, waitManifest, err := a.fileStore.Store(r, int64(len(manifest)), toEncrypt)
|
key, waitManifest, err := a.fileStore.Store(ctx, r, int64(len(manifest)), toEncrypt)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiPutFail.Inc(1)
|
apiPutFail.Inc(1)
|
||||||
return nil, nil, err
|
return nil, nil, err
|
||||||
|
|
@ -310,10 +310,10 @@ func (a *API) Put(content, contentType string, toEncrypt bool) (k storage.Addres
|
||||||
// Get uses iterative manifest retrieval and prefix matching
|
// Get uses iterative manifest retrieval and prefix matching
|
||||||
// to resolve basePath to content using FileStore retrieve
|
// to resolve basePath to content using FileStore retrieve
|
||||||
// it returns a section reader, mimeType, status, the key of the actual content and an error
|
// it returns a section reader, mimeType, status, the key of the actual content and an error
|
||||||
func (a *API) Get(manifestAddr storage.Address, path string) (reader storage.LazySectionReader, mimeType string, status int, contentAddr storage.Address, err error) {
|
func (a *API) Get(ctx context.Context, manifestAddr storage.Address, path string) (reader storage.LazySectionReader, mimeType string, status int, contentAddr storage.Address, err error) {
|
||||||
log.Debug("api.get", "key", manifestAddr, "path", path)
|
log.Debug("api.get", "key", manifestAddr, "path", path)
|
||||||
apiGetCount.Inc(1)
|
apiGetCount.Inc(1)
|
||||||
trie, err := loadManifest(a.fileStore, manifestAddr, nil)
|
trie, err := loadManifest(ctx, a.fileStore, manifestAddr, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiGetNotFound.Inc(1)
|
apiGetNotFound.Inc(1)
|
||||||
status = http.StatusNotFound
|
status = http.StatusNotFound
|
||||||
|
|
@ -375,7 +375,7 @@ func (a *API) Get(manifestAddr storage.Address, path string) (reader storage.Laz
|
||||||
log.Trace("resource is multihash", "key", manifestAddr)
|
log.Trace("resource is multihash", "key", manifestAddr)
|
||||||
|
|
||||||
// get the manifest the multihash digest points to
|
// get the manifest the multihash digest points to
|
||||||
trie, err := loadManifest(a.fileStore, manifestAddr, nil)
|
trie, err := loadManifest(ctx, a.fileStore, manifestAddr, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiGetNotFound.Inc(1)
|
apiGetNotFound.Inc(1)
|
||||||
status = http.StatusNotFound
|
status = http.StatusNotFound
|
||||||
|
|
@ -410,7 +410,7 @@ func (a *API) Get(manifestAddr storage.Address, path string) (reader storage.Laz
|
||||||
}
|
}
|
||||||
mimeType = entry.ContentType
|
mimeType = entry.ContentType
|
||||||
log.Debug("content lookup key", "key", contentAddr, "mimetype", mimeType)
|
log.Debug("content lookup key", "key", contentAddr, "mimetype", mimeType)
|
||||||
reader, _ = a.fileStore.Retrieve(contentAddr)
|
reader, _ = a.fileStore.Retrieve(ctx, contentAddr)
|
||||||
} else {
|
} else {
|
||||||
// no entry found
|
// no entry found
|
||||||
status = http.StatusNotFound
|
status = http.StatusNotFound
|
||||||
|
|
@ -422,10 +422,10 @@ func (a *API) Get(manifestAddr storage.Address, path string) (reader storage.Laz
|
||||||
}
|
}
|
||||||
|
|
||||||
// Modify loads manifest and checks the content hash before recalculating and storing the manifest.
|
// Modify loads manifest and checks the content hash before recalculating and storing the manifest.
|
||||||
func (a *API) Modify(addr storage.Address, path, contentHash, contentType string) (storage.Address, error) {
|
func (a *API) Modify(ctx context.Context, addr storage.Address, path, contentHash, contentType string) (storage.Address, error) {
|
||||||
apiModifyCount.Inc(1)
|
apiModifyCount.Inc(1)
|
||||||
quitC := make(chan bool)
|
quitC := make(chan bool)
|
||||||
trie, err := loadManifest(a.fileStore, addr, quitC)
|
trie, err := loadManifest(ctx, a.fileStore, addr, quitC)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiModifyFail.Inc(1)
|
apiModifyFail.Inc(1)
|
||||||
return nil, err
|
return nil, err
|
||||||
|
|
@ -449,7 +449,7 @@ func (a *API) Modify(addr storage.Address, path, contentHash, contentType string
|
||||||
}
|
}
|
||||||
|
|
||||||
// AddFile creates a new manifest entry, adds it to swarm, then adds a file to swarm.
|
// AddFile creates a new manifest entry, adds it to swarm, then adds a file to swarm.
|
||||||
func (a *API) AddFile(mhash, path, fname string, content []byte, nameresolver bool) (storage.Address, string, error) {
|
func (a *API) AddFile(ctx context.Context, mhash, path, fname string, content []byte, nameresolver bool) (storage.Address, string, error) {
|
||||||
apiAddFileCount.Inc(1)
|
apiAddFileCount.Inc(1)
|
||||||
|
|
||||||
uri, err := Parse("bzz:/" + mhash)
|
uri, err := Parse("bzz:/" + mhash)
|
||||||
|
|
@ -457,7 +457,7 @@ func (a *API) AddFile(mhash, path, fname string, content []byte, nameresolver bo
|
||||||
apiAddFileFail.Inc(1)
|
apiAddFileFail.Inc(1)
|
||||||
return nil, "", err
|
return nil, "", err
|
||||||
}
|
}
|
||||||
mkey, err := a.Resolve(uri)
|
mkey, err := a.Resolve(ctx, uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiAddFileFail.Inc(1)
|
apiAddFileFail.Inc(1)
|
||||||
return nil, "", err
|
return nil, "", err
|
||||||
|
|
@ -476,13 +476,13 @@ func (a *API) AddFile(mhash, path, fname string, content []byte, nameresolver bo
|
||||||
ModTime: time.Now(),
|
ModTime: time.Now(),
|
||||||
}
|
}
|
||||||
|
|
||||||
mw, err := a.NewManifestWriter(mkey, nil)
|
mw, err := a.NewManifestWriter(ctx, mkey, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiAddFileFail.Inc(1)
|
apiAddFileFail.Inc(1)
|
||||||
return nil, "", err
|
return nil, "", err
|
||||||
}
|
}
|
||||||
|
|
||||||
fkey, err := mw.AddEntry(bytes.NewReader(content), entry)
|
fkey, err := mw.AddEntry(ctx, bytes.NewReader(content), entry)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiAddFileFail.Inc(1)
|
apiAddFileFail.Inc(1)
|
||||||
return nil, "", err
|
return nil, "", err
|
||||||
|
|
@ -496,11 +496,10 @@ func (a *API) AddFile(mhash, path, fname string, content []byte, nameresolver bo
|
||||||
}
|
}
|
||||||
|
|
||||||
return fkey, newMkey.String(), nil
|
return fkey, newMkey.String(), nil
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// RemoveFile removes a file entry in a manifest.
|
// RemoveFile removes a file entry in a manifest.
|
||||||
func (a *API) RemoveFile(mhash, path, fname string, nameresolver bool) (string, error) {
|
func (a *API) RemoveFile(ctx context.Context, mhash string, path string, fname string, nameresolver bool) (string, error) {
|
||||||
apiRmFileCount.Inc(1)
|
apiRmFileCount.Inc(1)
|
||||||
|
|
||||||
uri, err := Parse("bzz:/" + mhash)
|
uri, err := Parse("bzz:/" + mhash)
|
||||||
|
|
@ -508,7 +507,7 @@ func (a *API) RemoveFile(mhash, path, fname string, nameresolver bool) (string,
|
||||||
apiRmFileFail.Inc(1)
|
apiRmFileFail.Inc(1)
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
mkey, err := a.Resolve(uri)
|
mkey, err := a.Resolve(ctx, uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiRmFileFail.Inc(1)
|
apiRmFileFail.Inc(1)
|
||||||
return "", err
|
return "", err
|
||||||
|
|
@ -519,7 +518,7 @@ func (a *API) RemoveFile(mhash, path, fname string, nameresolver bool) (string,
|
||||||
path = path[1:]
|
path = path[1:]
|
||||||
}
|
}
|
||||||
|
|
||||||
mw, err := a.NewManifestWriter(mkey, nil)
|
mw, err := a.NewManifestWriter(ctx, mkey, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiRmFileFail.Inc(1)
|
apiRmFileFail.Inc(1)
|
||||||
return "", err
|
return "", err
|
||||||
|
|
@ -542,7 +541,7 @@ func (a *API) RemoveFile(mhash, path, fname string, nameresolver bool) (string,
|
||||||
}
|
}
|
||||||
|
|
||||||
// AppendFile removes old manifest, appends file entry to new manifest and adds it to Swarm.
|
// AppendFile removes old manifest, appends file entry to new manifest and adds it to Swarm.
|
||||||
func (a *API) AppendFile(mhash, path, fname string, existingSize int64, content []byte, oldAddr storage.Address, offset int64, addSize int64, nameresolver bool) (storage.Address, string, error) {
|
func (a *API) AppendFile(ctx context.Context, mhash, path, fname string, existingSize int64, content []byte, oldAddr storage.Address, offset int64, addSize int64, nameresolver bool) (storage.Address, string, error) {
|
||||||
apiAppendFileCount.Inc(1)
|
apiAppendFileCount.Inc(1)
|
||||||
|
|
||||||
buffSize := offset + addSize
|
buffSize := offset + addSize
|
||||||
|
|
@ -552,7 +551,7 @@ func (a *API) AppendFile(mhash, path, fname string, existingSize int64, content
|
||||||
|
|
||||||
buf := make([]byte, buffSize)
|
buf := make([]byte, buffSize)
|
||||||
|
|
||||||
oldReader, _ := a.Retrieve(oldAddr)
|
oldReader, _ := a.Retrieve(ctx, oldAddr)
|
||||||
io.ReadAtLeast(oldReader, buf, int(offset))
|
io.ReadAtLeast(oldReader, buf, int(offset))
|
||||||
|
|
||||||
newReader := bytes.NewReader(content)
|
newReader := bytes.NewReader(content)
|
||||||
|
|
@ -575,7 +574,7 @@ func (a *API) AppendFile(mhash, path, fname string, existingSize int64, content
|
||||||
apiAppendFileFail.Inc(1)
|
apiAppendFileFail.Inc(1)
|
||||||
return nil, "", err
|
return nil, "", err
|
||||||
}
|
}
|
||||||
mkey, err := a.Resolve(uri)
|
mkey, err := a.Resolve(ctx, uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiAppendFileFail.Inc(1)
|
apiAppendFileFail.Inc(1)
|
||||||
return nil, "", err
|
return nil, "", err
|
||||||
|
|
@ -586,7 +585,7 @@ func (a *API) AppendFile(mhash, path, fname string, existingSize int64, content
|
||||||
path = path[1:]
|
path = path[1:]
|
||||||
}
|
}
|
||||||
|
|
||||||
mw, err := a.NewManifestWriter(mkey, nil)
|
mw, err := a.NewManifestWriter(ctx, mkey, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiAppendFileFail.Inc(1)
|
apiAppendFileFail.Inc(1)
|
||||||
return nil, "", err
|
return nil, "", err
|
||||||
|
|
@ -606,7 +605,7 @@ func (a *API) AppendFile(mhash, path, fname string, existingSize int64, content
|
||||||
ModTime: time.Now(),
|
ModTime: time.Now(),
|
||||||
}
|
}
|
||||||
|
|
||||||
fkey, err := mw.AddEntry(io.Reader(combinedReader), entry)
|
fkey, err := mw.AddEntry(ctx, io.Reader(combinedReader), entry)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
apiAppendFileFail.Inc(1)
|
apiAppendFileFail.Inc(1)
|
||||||
return nil, "", err
|
return nil, "", err
|
||||||
|
|
@ -620,23 +619,22 @@ func (a *API) AppendFile(mhash, path, fname string, existingSize int64, content
|
||||||
}
|
}
|
||||||
|
|
||||||
return fkey, newMkey.String(), nil
|
return fkey, newMkey.String(), nil
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// BuildDirectoryTree used by swarmfs_unix
|
// BuildDirectoryTree used by swarmfs_unix
|
||||||
func (a *API) BuildDirectoryTree(mhash string, nameresolver bool) (addr storage.Address, manifestEntryMap map[string]*manifestTrieEntry, err error) {
|
func (a *API) BuildDirectoryTree(ctx context.Context, mhash string, nameresolver bool) (addr storage.Address, manifestEntryMap map[string]*manifestTrieEntry, err error) {
|
||||||
|
|
||||||
uri, err := Parse("bzz:/" + mhash)
|
uri, err := Parse("bzz:/" + mhash)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, nil, err
|
return nil, nil, err
|
||||||
}
|
}
|
||||||
addr, err = a.Resolve(uri)
|
addr, err = a.Resolve(ctx, uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, nil, err
|
return nil, nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
quitC := make(chan bool)
|
quitC := make(chan bool)
|
||||||
rootTrie, err := loadManifest(a.fileStore, addr, quitC)
|
rootTrie, err := loadManifest(ctx, a.fileStore, addr, quitC)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, nil, fmt.Errorf("can't load manifest %v: %v", addr.String(), err)
|
return nil, nil, fmt.Errorf("can't load manifest %v: %v", addr.String(), err)
|
||||||
}
|
}
|
||||||
|
|
@ -725,8 +723,8 @@ func (a *API) ResourceIsValidated() bool {
|
||||||
}
|
}
|
||||||
|
|
||||||
// ResolveResourceManifest retrieves the Mutable Resource manifest for the given address, and returns the address of the metadata chunk.
|
// ResolveResourceManifest retrieves the Mutable Resource manifest for the given address, and returns the address of the metadata chunk.
|
||||||
func (a *API) ResolveResourceManifest(addr storage.Address) (storage.Address, error) {
|
func (a *API) ResolveResourceManifest(ctx context.Context, addr storage.Address) (storage.Address, error) {
|
||||||
trie, err := loadManifest(a.fileStore, addr, nil)
|
trie, err := loadManifest(ctx, a.fileStore, addr, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("cannot load resource manifest: %v", err)
|
return nil, fmt.Errorf("cannot load resource manifest: %v", err)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -85,7 +85,7 @@ func expResponse(content string, mimeType string, status int) *Response {
|
||||||
|
|
||||||
func testGet(t *testing.T, api *API, bzzhash, path string) *testResponse {
|
func testGet(t *testing.T, api *API, bzzhash, path string) *testResponse {
|
||||||
addr := storage.Address(common.Hex2Bytes(bzzhash))
|
addr := storage.Address(common.Hex2Bytes(bzzhash))
|
||||||
reader, mimeType, status, _, err := api.Get(addr, path)
|
reader, mimeType, status, _, err := api.Get(context.TODO(), addr, path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("unexpected error: %v", err)
|
t.Fatalf("unexpected error: %v", err)
|
||||||
}
|
}
|
||||||
|
|
@ -109,8 +109,7 @@ func TestApiPut(t *testing.T) {
|
||||||
testAPI(t, func(api *API, toEncrypt bool) {
|
testAPI(t, func(api *API, toEncrypt bool) {
|
||||||
content := "hello"
|
content := "hello"
|
||||||
exp := expResponse(content, "text/plain", 0)
|
exp := expResponse(content, "text/plain", 0)
|
||||||
// exp := expResponse([]byte(content), "text/plain", 0)
|
addr, wait, err := api.Put(context.TODO(), content, exp.MimeType, toEncrypt)
|
||||||
addr, wait, err := api.Put(content, exp.MimeType, toEncrypt)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("unexpected error: %v", err)
|
t.Fatalf("unexpected error: %v", err)
|
||||||
}
|
}
|
||||||
|
|
@ -226,7 +225,7 @@ func TestAPIResolve(t *testing.T) {
|
||||||
if x.immutable {
|
if x.immutable {
|
||||||
uri.Scheme = "bzz-immutable"
|
uri.Scheme = "bzz-immutable"
|
||||||
}
|
}
|
||||||
res, err := api.Resolve(uri)
|
res, err := api.Resolve(context.TODO(), uri)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
if x.expectErr != nil {
|
if x.expectErr != nil {
|
||||||
t.Fatalf("expected error %q, got result %q", x.expectErr, res)
|
t.Fatalf("expected error %q, got result %q", x.expectErr, res)
|
||||||
|
|
|
||||||
|
|
@ -18,6 +18,7 @@ package api
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bufio"
|
"bufio"
|
||||||
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
|
@ -114,7 +115,7 @@ func (fs *FileSystem) Upload(lpath, index string, toEncrypt bool) (string, error
|
||||||
stat, _ := f.Stat()
|
stat, _ := f.Stat()
|
||||||
var hash storage.Address
|
var hash storage.Address
|
||||||
var wait func()
|
var wait func()
|
||||||
hash, wait, err = fs.api.fileStore.Store(f, stat.Size(), toEncrypt)
|
hash, wait, err = fs.api.fileStore.Store(context.TODO(), f, stat.Size(), toEncrypt)
|
||||||
if hash != nil {
|
if hash != nil {
|
||||||
list[i].Hash = hash.Hex()
|
list[i].Hash = hash.Hex()
|
||||||
}
|
}
|
||||||
|
|
@ -189,7 +190,7 @@ func (fs *FileSystem) Download(bzzpath, localpath string) error {
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
addr, err := fs.api.Resolve(uri)
|
addr, err := fs.api.Resolve(context.TODO(), uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
@ -200,7 +201,7 @@ func (fs *FileSystem) Download(bzzpath, localpath string) error {
|
||||||
}
|
}
|
||||||
|
|
||||||
quitC := make(chan bool)
|
quitC := make(chan bool)
|
||||||
trie, err := loadManifest(fs.api.fileStore, addr, quitC)
|
trie, err := loadManifest(context.TODO(), fs.api.fileStore, addr, quitC)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Warn(fmt.Sprintf("fs.Download: loadManifestTrie error: %v", err))
|
log.Warn(fmt.Sprintf("fs.Download: loadManifestTrie error: %v", err))
|
||||||
return err
|
return err
|
||||||
|
|
@ -273,7 +274,7 @@ func retrieveToFile(quitC chan bool, fileStore *storage.FileStore, addr storage.
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
reader, _ := fileStore.Retrieve(addr)
|
reader, _ := fileStore.Retrieve(context.TODO(), addr)
|
||||||
writer := bufio.NewWriter(f)
|
writer := bufio.NewWriter(f)
|
||||||
size, err := reader.Size(quitC)
|
size, err := reader.Size(quitC)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
|
||||||
|
|
@ -18,6 +18,7 @@ package api
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"context"
|
||||||
"io/ioutil"
|
"io/ioutil"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
|
@ -63,7 +64,7 @@ func TestApiDirUpload0(t *testing.T) {
|
||||||
checkResponse(t, resp, exp)
|
checkResponse(t, resp, exp)
|
||||||
|
|
||||||
addr := storage.Address(common.Hex2Bytes(bzzhash))
|
addr := storage.Address(common.Hex2Bytes(bzzhash))
|
||||||
_, _, _, _, err = api.Get(addr, "")
|
_, _, _, _, err = api.Get(context.TODO(), addr, "")
|
||||||
if err == nil {
|
if err == nil {
|
||||||
t.Fatalf("expected error: %v", err)
|
t.Fatalf("expected error: %v", err)
|
||||||
}
|
}
|
||||||
|
|
@ -95,7 +96,7 @@ func TestApiDirUploadModify(t *testing.T) {
|
||||||
}
|
}
|
||||||
|
|
||||||
addr := storage.Address(common.Hex2Bytes(bzzhash))
|
addr := storage.Address(common.Hex2Bytes(bzzhash))
|
||||||
addr, err = api.Modify(addr, "index.html", "", "")
|
addr, err = api.Modify(context.TODO(), addr, "index.html", "", "")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Errorf("unexpected error: %v", err)
|
t.Errorf("unexpected error: %v", err)
|
||||||
return
|
return
|
||||||
|
|
@ -105,18 +106,18 @@ func TestApiDirUploadModify(t *testing.T) {
|
||||||
t.Errorf("unexpected error: %v", err)
|
t.Errorf("unexpected error: %v", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
hash, wait, err := api.Store(bytes.NewReader(index), int64(len(index)), toEncrypt)
|
hash, wait, err := api.Store(context.TODO(), bytes.NewReader(index), int64(len(index)), toEncrypt)
|
||||||
wait()
|
wait()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Errorf("unexpected error: %v", err)
|
t.Errorf("unexpected error: %v", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
addr, err = api.Modify(addr, "index2.html", hash.Hex(), "text/html; charset=utf-8")
|
addr, err = api.Modify(context.TODO(), addr, "index2.html", hash.Hex(), "text/html; charset=utf-8")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Errorf("unexpected error: %v", err)
|
t.Errorf("unexpected error: %v", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
addr, err = api.Modify(addr, "img/logo.png", hash.Hex(), "text/html; charset=utf-8")
|
addr, err = api.Modify(context.TODO(), addr, "img/logo.png", hash.Hex(), "text/html; charset=utf-8")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Errorf("unexpected error: %v", err)
|
t.Errorf("unexpected error: %v", err)
|
||||||
return
|
return
|
||||||
|
|
@ -137,7 +138,7 @@ func TestApiDirUploadModify(t *testing.T) {
|
||||||
exp = expResponse(content, "text/css", 0)
|
exp = expResponse(content, "text/css", 0)
|
||||||
checkResponse(t, resp, exp)
|
checkResponse(t, resp, exp)
|
||||||
|
|
||||||
_, _, _, _, err = api.Get(addr, "")
|
_, _, _, _, err = api.Get(context.TODO(), addr, "")
|
||||||
if err == nil {
|
if err == nil {
|
||||||
t.Errorf("expected error: %v", err)
|
t.Errorf("expected error: %v", err)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -23,6 +23,7 @@ import (
|
||||||
"archive/tar"
|
"archive/tar"
|
||||||
"bufio"
|
"bufio"
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
|
@ -120,7 +121,7 @@ type Request struct {
|
||||||
|
|
||||||
// HandlePostRaw handles a POST request to a raw bzz-raw:/ URI, stores the request
|
// HandlePostRaw handles a POST request to a raw bzz-raw:/ URI, stores the request
|
||||||
// body in swarm and returns the resulting storage address as a text/plain response
|
// body in swarm and returns the resulting storage address as a text/plain response
|
||||||
func (s *Server) HandlePostRaw(w http.ResponseWriter, r *Request) {
|
func (s *Server) HandlePostRaw(ctx context.Context, w http.ResponseWriter, r *Request) {
|
||||||
log.Debug("handle.post.raw", "ruid", r.ruid)
|
log.Debug("handle.post.raw", "ruid", r.ruid)
|
||||||
|
|
||||||
postRawCount.Inc(1)
|
postRawCount.Inc(1)
|
||||||
|
|
@ -147,7 +148,7 @@ func (s *Server) HandlePostRaw(w http.ResponseWriter, r *Request) {
|
||||||
Respond(w, r, "missing Content-Length header in request", http.StatusBadRequest)
|
Respond(w, r, "missing Content-Length header in request", http.StatusBadRequest)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
addr, _, err := s.api.Store(r.Body, r.ContentLength, toEncrypt)
|
addr, _, err := s.api.Store(ctx, r.Body, r.ContentLength, toEncrypt)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
postRawFail.Inc(1)
|
postRawFail.Inc(1)
|
||||||
Respond(w, r, err.Error(), http.StatusInternalServerError)
|
Respond(w, r, err.Error(), http.StatusInternalServerError)
|
||||||
|
|
@ -166,7 +167,7 @@ func (s *Server) HandlePostRaw(w http.ResponseWriter, r *Request) {
|
||||||
// (either a tar archive or multipart form), adds those files either to an
|
// (either a tar archive or multipart form), adds those files either to an
|
||||||
// existing manifest or to a new manifest under <path> and returns the
|
// existing manifest or to a new manifest under <path> and returns the
|
||||||
// resulting manifest hash as a text/plain response
|
// resulting manifest hash as a text/plain response
|
||||||
func (s *Server) HandlePostFiles(w http.ResponseWriter, r *Request) {
|
func (s *Server) HandlePostFiles(ctx context.Context, w http.ResponseWriter, r *Request) {
|
||||||
log.Debug("handle.post.files", "ruid", r.ruid)
|
log.Debug("handle.post.files", "ruid", r.ruid)
|
||||||
|
|
||||||
postFilesCount.Inc(1)
|
postFilesCount.Inc(1)
|
||||||
|
|
@ -184,7 +185,7 @@ func (s *Server) HandlePostFiles(w http.ResponseWriter, r *Request) {
|
||||||
|
|
||||||
var addr storage.Address
|
var addr storage.Address
|
||||||
if r.uri.Addr != "" && r.uri.Addr != "encrypt" {
|
if r.uri.Addr != "" && r.uri.Addr != "encrypt" {
|
||||||
addr, err = s.api.Resolve(r.uri)
|
addr, err = s.api.Resolve(ctx, r.uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
postFilesFail.Inc(1)
|
postFilesFail.Inc(1)
|
||||||
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusInternalServerError)
|
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusInternalServerError)
|
||||||
|
|
@ -192,7 +193,7 @@ func (s *Server) HandlePostFiles(w http.ResponseWriter, r *Request) {
|
||||||
}
|
}
|
||||||
log.Debug("resolved key", "ruid", r.ruid, "key", addr)
|
log.Debug("resolved key", "ruid", r.ruid, "key", addr)
|
||||||
} else {
|
} else {
|
||||||
addr, err = s.api.NewManifest(toEncrypt)
|
addr, err = s.api.NewManifest(ctx, toEncrypt)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
postFilesFail.Inc(1)
|
postFilesFail.Inc(1)
|
||||||
Respond(w, r, err.Error(), http.StatusInternalServerError)
|
Respond(w, r, err.Error(), http.StatusInternalServerError)
|
||||||
|
|
@ -201,17 +202,17 @@ func (s *Server) HandlePostFiles(w http.ResponseWriter, r *Request) {
|
||||||
log.Debug("new manifest", "ruid", r.ruid, "key", addr)
|
log.Debug("new manifest", "ruid", r.ruid, "key", addr)
|
||||||
}
|
}
|
||||||
|
|
||||||
newAddr, err := s.updateManifest(addr, func(mw *api.ManifestWriter) error {
|
newAddr, err := s.updateManifest(ctx, addr, func(mw *api.ManifestWriter) error {
|
||||||
switch contentType {
|
switch contentType {
|
||||||
|
|
||||||
case "application/x-tar":
|
case "application/x-tar":
|
||||||
return s.handleTarUpload(r, mw)
|
return s.handleTarUpload(ctx, r, mw)
|
||||||
|
|
||||||
case "multipart/form-data":
|
case "multipart/form-data":
|
||||||
return s.handleMultipartUpload(r, params["boundary"], mw)
|
return s.handleMultipartUpload(ctx, r, params["boundary"], mw)
|
||||||
|
|
||||||
default:
|
default:
|
||||||
return s.handleDirectUpload(r, mw)
|
return s.handleDirectUpload(ctx, r, mw)
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
@ -227,7 +228,7 @@ func (s *Server) HandlePostFiles(w http.ResponseWriter, r *Request) {
|
||||||
fmt.Fprint(w, newAddr)
|
fmt.Fprint(w, newAddr)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) handleTarUpload(req *Request, mw *api.ManifestWriter) error {
|
func (s *Server) handleTarUpload(ctx context.Context, req *Request, mw *api.ManifestWriter) error {
|
||||||
log.Debug("handle.tar.upload", "ruid", req.ruid)
|
log.Debug("handle.tar.upload", "ruid", req.ruid)
|
||||||
tr := tar.NewReader(req.Body)
|
tr := tar.NewReader(req.Body)
|
||||||
for {
|
for {
|
||||||
|
|
@ -253,7 +254,7 @@ func (s *Server) handleTarUpload(req *Request, mw *api.ManifestWriter) error {
|
||||||
ModTime: hdr.ModTime,
|
ModTime: hdr.ModTime,
|
||||||
}
|
}
|
||||||
log.Debug("adding path to new manifest", "ruid", req.ruid, "bytes", entry.Size, "path", entry.Path)
|
log.Debug("adding path to new manifest", "ruid", req.ruid, "bytes", entry.Size, "path", entry.Path)
|
||||||
contentKey, err := mw.AddEntry(tr, entry)
|
contentKey, err := mw.AddEntry(ctx, tr, entry)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("error adding manifest entry from tar stream: %s", err)
|
return fmt.Errorf("error adding manifest entry from tar stream: %s", err)
|
||||||
}
|
}
|
||||||
|
|
@ -261,7 +262,7 @@ func (s *Server) handleTarUpload(req *Request, mw *api.ManifestWriter) error {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) handleMultipartUpload(req *Request, boundary string, mw *api.ManifestWriter) error {
|
func (s *Server) handleMultipartUpload(ctx context.Context, req *Request, boundary string, mw *api.ManifestWriter) error {
|
||||||
log.Debug("handle.multipart.upload", "ruid", req.ruid)
|
log.Debug("handle.multipart.upload", "ruid", req.ruid)
|
||||||
mr := multipart.NewReader(req.Body, boundary)
|
mr := multipart.NewReader(req.Body, boundary)
|
||||||
for {
|
for {
|
||||||
|
|
@ -311,7 +312,7 @@ func (s *Server) handleMultipartUpload(req *Request, boundary string, mw *api.Ma
|
||||||
ModTime: time.Now(),
|
ModTime: time.Now(),
|
||||||
}
|
}
|
||||||
log.Debug("adding path to new manifest", "ruid", req.ruid, "bytes", entry.Size, "path", entry.Path)
|
log.Debug("adding path to new manifest", "ruid", req.ruid, "bytes", entry.Size, "path", entry.Path)
|
||||||
contentKey, err := mw.AddEntry(reader, entry)
|
contentKey, err := mw.AddEntry(ctx, reader, entry)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("error adding manifest entry from multipart form: %s", err)
|
return fmt.Errorf("error adding manifest entry from multipart form: %s", err)
|
||||||
}
|
}
|
||||||
|
|
@ -319,9 +320,9 @@ func (s *Server) handleMultipartUpload(req *Request, boundary string, mw *api.Ma
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) handleDirectUpload(req *Request, mw *api.ManifestWriter) error {
|
func (s *Server) handleDirectUpload(ctx context.Context, req *Request, mw *api.ManifestWriter) error {
|
||||||
log.Debug("handle.direct.upload", "ruid", req.ruid)
|
log.Debug("handle.direct.upload", "ruid", req.ruid)
|
||||||
key, err := mw.AddEntry(req.Body, &api.ManifestEntry{
|
key, err := mw.AddEntry(ctx, req.Body, &api.ManifestEntry{
|
||||||
Path: req.uri.Path,
|
Path: req.uri.Path,
|
||||||
ContentType: req.Header.Get("Content-Type"),
|
ContentType: req.Header.Get("Content-Type"),
|
||||||
Mode: 0644,
|
Mode: 0644,
|
||||||
|
|
@ -338,18 +339,18 @@ func (s *Server) handleDirectUpload(req *Request, mw *api.ManifestWriter) error
|
||||||
// HandleDelete handles a DELETE request to bzz:/<manifest>/<path>, removes
|
// HandleDelete handles a DELETE request to bzz:/<manifest>/<path>, removes
|
||||||
// <path> from <manifest> and returns the resulting manifest hash as a
|
// <path> from <manifest> and returns the resulting manifest hash as a
|
||||||
// text/plain response
|
// text/plain response
|
||||||
func (s *Server) HandleDelete(w http.ResponseWriter, r *Request) {
|
func (s *Server) HandleDelete(ctx context.Context, w http.ResponseWriter, r *Request) {
|
||||||
log.Debug("handle.delete", "ruid", r.ruid)
|
log.Debug("handle.delete", "ruid", r.ruid)
|
||||||
|
|
||||||
deleteCount.Inc(1)
|
deleteCount.Inc(1)
|
||||||
key, err := s.api.Resolve(r.uri)
|
key, err := s.api.Resolve(ctx, r.uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
deleteFail.Inc(1)
|
deleteFail.Inc(1)
|
||||||
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusInternalServerError)
|
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusInternalServerError)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
newKey, err := s.updateManifest(key, func(mw *api.ManifestWriter) error {
|
newKey, err := s.updateManifest(ctx, key, func(mw *api.ManifestWriter) error {
|
||||||
log.Debug(fmt.Sprintf("removing %s from manifest %s", r.uri.Path, key.Log()), "ruid", r.ruid)
|
log.Debug(fmt.Sprintf("removing %s from manifest %s", r.uri.Path, key.Log()), "ruid", r.ruid)
|
||||||
return mw.RemoveEntry(r.uri.Path)
|
return mw.RemoveEntry(r.uri.Path)
|
||||||
})
|
})
|
||||||
|
|
@ -399,7 +400,7 @@ func resourcePostMode(path string) (isRaw bool, frequency uint64, err error) {
|
||||||
// The resource name will be verbatim what is passed as the address part of the url.
|
// The resource name will be verbatim what is passed as the address part of the url.
|
||||||
// For example, if a POST is made to /bzz-resource:/foo.eth/raw/13 a new resource with frequency 13
|
// For example, if a POST is made to /bzz-resource:/foo.eth/raw/13 a new resource with frequency 13
|
||||||
// and name "foo.eth" will be created
|
// and name "foo.eth" will be created
|
||||||
func (s *Server) HandlePostResource(w http.ResponseWriter, r *Request) {
|
func (s *Server) HandlePostResource(ctx context.Context, w http.ResponseWriter, r *Request) {
|
||||||
log.Debug("handle.post.resource", "ruid", r.ruid)
|
log.Debug("handle.post.resource", "ruid", r.ruid)
|
||||||
var err error
|
var err error
|
||||||
var addr storage.Address
|
var addr storage.Address
|
||||||
|
|
@ -428,7 +429,7 @@ func (s *Server) HandlePostResource(w http.ResponseWriter, r *Request) {
|
||||||
// we create a manifest so we can retrieve the resource with bzz:// later
|
// we create a manifest so we can retrieve the resource with bzz:// later
|
||||||
// this manifest has a special "resource type" manifest, and its hash is the key of the mutable resource
|
// this manifest has a special "resource type" manifest, and its hash is the key of the mutable resource
|
||||||
// root chunk
|
// root chunk
|
||||||
m, err := s.api.NewResourceManifest(addr.Hex())
|
m, err := s.api.NewResourceManifest(ctx, addr.Hex())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
Respond(w, r, fmt.Sprintf("failed to create resource manifest: %v", err), http.StatusInternalServerError)
|
Respond(w, r, fmt.Sprintf("failed to create resource manifest: %v", err), http.StatusInternalServerError)
|
||||||
return
|
return
|
||||||
|
|
@ -448,7 +449,7 @@ func (s *Server) HandlePostResource(w http.ResponseWriter, r *Request) {
|
||||||
// that means that we retrieve the manifest and inspect its Hash member.
|
// that means that we retrieve the manifest and inspect its Hash member.
|
||||||
manifestAddr := r.uri.Address()
|
manifestAddr := r.uri.Address()
|
||||||
if manifestAddr == nil {
|
if manifestAddr == nil {
|
||||||
manifestAddr, err = s.api.Resolve(r.uri)
|
manifestAddr, err = s.api.Resolve(ctx, r.uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
getFail.Inc(1)
|
getFail.Inc(1)
|
||||||
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
||||||
|
|
@ -459,7 +460,7 @@ func (s *Server) HandlePostResource(w http.ResponseWriter, r *Request) {
|
||||||
}
|
}
|
||||||
|
|
||||||
// get the root chunk key from the manifest
|
// get the root chunk key from the manifest
|
||||||
addr, err = s.api.ResolveResourceManifest(manifestAddr)
|
addr, err = s.api.ResolveResourceManifest(ctx, manifestAddr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
getFail.Inc(1)
|
getFail.Inc(1)
|
||||||
Respond(w, r, fmt.Sprintf("error resolving resource root chunk for %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
Respond(w, r, fmt.Sprintf("error resolving resource root chunk for %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
||||||
|
|
@ -518,19 +519,19 @@ func (s *Server) HandlePostResource(w http.ResponseWriter, r *Request) {
|
||||||
// bzz-resource://<id>/<n> - get latest update on period n
|
// bzz-resource://<id>/<n> - get latest update on period n
|
||||||
// bzz-resource://<id>/<n>/<m> - get update version m of period n
|
// bzz-resource://<id>/<n>/<m> - get update version m of period n
|
||||||
// <id> = ens name or hash
|
// <id> = ens name or hash
|
||||||
func (s *Server) HandleGetResource(w http.ResponseWriter, r *Request) {
|
func (s *Server) HandleGetResource(ctx context.Context, w http.ResponseWriter, r *Request) {
|
||||||
s.handleGetResource(w, r)
|
s.handleGetResource(ctx, w, r)
|
||||||
}
|
}
|
||||||
|
|
||||||
// TODO: Enable pass maxPeriod parameter
|
// TODO: Enable pass maxPeriod parameter
|
||||||
func (s *Server) handleGetResource(w http.ResponseWriter, r *Request) {
|
func (s *Server) handleGetResource(ctx context.Context, w http.ResponseWriter, r *Request) {
|
||||||
log.Debug("handle.get.resource", "ruid", r.ruid)
|
log.Debug("handle.get.resource", "ruid", r.ruid)
|
||||||
var err error
|
var err error
|
||||||
|
|
||||||
// resolve the content key.
|
// resolve the content key.
|
||||||
manifestAddr := r.uri.Address()
|
manifestAddr := r.uri.Address()
|
||||||
if manifestAddr == nil {
|
if manifestAddr == nil {
|
||||||
manifestAddr, err = s.api.Resolve(r.uri)
|
manifestAddr, err = s.api.Resolve(ctx, r.uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
getFail.Inc(1)
|
getFail.Inc(1)
|
||||||
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
||||||
|
|
@ -541,7 +542,7 @@ func (s *Server) handleGetResource(w http.ResponseWriter, r *Request) {
|
||||||
}
|
}
|
||||||
|
|
||||||
// get the root chunk key from the manifest
|
// get the root chunk key from the manifest
|
||||||
key, err := s.api.ResolveResourceManifest(manifestAddr)
|
key, err := s.api.ResolveResourceManifest(ctx, manifestAddr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
getFail.Inc(1)
|
getFail.Inc(1)
|
||||||
Respond(w, r, fmt.Sprintf("error resolving resource root chunk for %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
Respond(w, r, fmt.Sprintf("error resolving resource root chunk for %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
||||||
|
|
@ -623,13 +624,13 @@ func (s *Server) translateResourceError(w http.ResponseWriter, r *Request, supEr
|
||||||
// given storage key
|
// given storage key
|
||||||
// - bzz-hash://<key> and responds with the hash of the content stored
|
// - bzz-hash://<key> and responds with the hash of the content stored
|
||||||
// at the given storage key as a text/plain response
|
// at the given storage key as a text/plain response
|
||||||
func (s *Server) HandleGet(w http.ResponseWriter, r *Request) {
|
func (s *Server) HandleGet(ctx context.Context, w http.ResponseWriter, r *Request) {
|
||||||
log.Debug("handle.get", "ruid", r.ruid, "uri", r.uri)
|
log.Debug("handle.get", "ruid", r.ruid, "uri", r.uri)
|
||||||
getCount.Inc(1)
|
getCount.Inc(1)
|
||||||
var err error
|
var err error
|
||||||
addr := r.uri.Address()
|
addr := r.uri.Address()
|
||||||
if addr == nil {
|
if addr == nil {
|
||||||
addr, err = s.api.Resolve(r.uri)
|
addr, err = s.api.Resolve(ctx, r.uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
getFail.Inc(1)
|
getFail.Inc(1)
|
||||||
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
||||||
|
|
@ -644,7 +645,7 @@ func (s *Server) HandleGet(w http.ResponseWriter, r *Request) {
|
||||||
// if path is set, interpret <key> as a manifest and return the
|
// if path is set, interpret <key> as a manifest and return the
|
||||||
// raw entry at the given path
|
// raw entry at the given path
|
||||||
if r.uri.Path != "" {
|
if r.uri.Path != "" {
|
||||||
walker, err := s.api.NewManifestWalker(addr, nil)
|
walker, err := s.api.NewManifestWalker(ctx, addr, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
getFail.Inc(1)
|
getFail.Inc(1)
|
||||||
Respond(w, r, fmt.Sprintf("%s is not a manifest", addr), http.StatusBadRequest)
|
Respond(w, r, fmt.Sprintf("%s is not a manifest", addr), http.StatusBadRequest)
|
||||||
|
|
@ -692,7 +693,7 @@ func (s *Server) HandleGet(w http.ResponseWriter, r *Request) {
|
||||||
}
|
}
|
||||||
|
|
||||||
// check the root chunk exists by retrieving the file's size
|
// check the root chunk exists by retrieving the file's size
|
||||||
reader, isEncrypted := s.api.Retrieve(addr)
|
reader, isEncrypted := s.api.Retrieve(ctx, addr)
|
||||||
if _, err := reader.Size(nil); err != nil {
|
if _, err := reader.Size(nil); err != nil {
|
||||||
getFail.Inc(1)
|
getFail.Inc(1)
|
||||||
Respond(w, r, fmt.Sprintf("root chunk not found %s: %s", addr, err), http.StatusNotFound)
|
Respond(w, r, fmt.Sprintf("root chunk not found %s: %s", addr, err), http.StatusNotFound)
|
||||||
|
|
@ -721,7 +722,7 @@ func (s *Server) HandleGet(w http.ResponseWriter, r *Request) {
|
||||||
// HandleGetFiles handles a GET request to bzz:/<manifest> with an Accept
|
// HandleGetFiles handles a GET request to bzz:/<manifest> with an Accept
|
||||||
// header of "application/x-tar" and returns a tar stream of all files
|
// header of "application/x-tar" and returns a tar stream of all files
|
||||||
// contained in the manifest
|
// contained in the manifest
|
||||||
func (s *Server) HandleGetFiles(w http.ResponseWriter, r *Request) {
|
func (s *Server) HandleGetFiles(ctx context.Context, w http.ResponseWriter, r *Request) {
|
||||||
log.Debug("handle.get.files", "ruid", r.ruid, "uri", r.uri)
|
log.Debug("handle.get.files", "ruid", r.ruid, "uri", r.uri)
|
||||||
getFilesCount.Inc(1)
|
getFilesCount.Inc(1)
|
||||||
if r.uri.Path != "" {
|
if r.uri.Path != "" {
|
||||||
|
|
@ -730,7 +731,7 @@ func (s *Server) HandleGetFiles(w http.ResponseWriter, r *Request) {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
addr, err := s.api.Resolve(r.uri)
|
addr, err := s.api.Resolve(ctx, r.uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
getFilesFail.Inc(1)
|
getFilesFail.Inc(1)
|
||||||
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
||||||
|
|
@ -738,7 +739,7 @@ func (s *Server) HandleGetFiles(w http.ResponseWriter, r *Request) {
|
||||||
}
|
}
|
||||||
log.Debug("handle.get.files: resolved", "ruid", r.ruid, "key", addr)
|
log.Debug("handle.get.files: resolved", "ruid", r.ruid, "key", addr)
|
||||||
|
|
||||||
walker, err := s.api.NewManifestWalker(addr, nil)
|
walker, err := s.api.NewManifestWalker(ctx, addr, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
getFilesFail.Inc(1)
|
getFilesFail.Inc(1)
|
||||||
Respond(w, r, err.Error(), http.StatusInternalServerError)
|
Respond(w, r, err.Error(), http.StatusInternalServerError)
|
||||||
|
|
@ -757,7 +758,7 @@ func (s *Server) HandleGetFiles(w http.ResponseWriter, r *Request) {
|
||||||
}
|
}
|
||||||
|
|
||||||
// retrieve the entry's key and size
|
// retrieve the entry's key and size
|
||||||
reader, isEncrypted := s.api.Retrieve(storage.Address(common.Hex2Bytes(entry.Hash)))
|
reader, isEncrypted := s.api.Retrieve(ctx, storage.Address(common.Hex2Bytes(entry.Hash)))
|
||||||
size, err := reader.Size(nil)
|
size, err := reader.Size(nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
|
|
@ -797,7 +798,7 @@ func (s *Server) HandleGetFiles(w http.ResponseWriter, r *Request) {
|
||||||
// HandleGetList handles a GET request to bzz-list:/<manifest>/<path> and returns
|
// HandleGetList handles a GET request to bzz-list:/<manifest>/<path> and returns
|
||||||
// a list of all files contained in <manifest> under <path> grouped into
|
// a list of all files contained in <manifest> under <path> grouped into
|
||||||
// common prefixes using "/" as a delimiter
|
// common prefixes using "/" as a delimiter
|
||||||
func (s *Server) HandleGetList(w http.ResponseWriter, r *Request) {
|
func (s *Server) HandleGetList(ctx context.Context, w http.ResponseWriter, r *Request) {
|
||||||
log.Debug("handle.get.list", "ruid", r.ruid, "uri", r.uri)
|
log.Debug("handle.get.list", "ruid", r.ruid, "uri", r.uri)
|
||||||
getListCount.Inc(1)
|
getListCount.Inc(1)
|
||||||
// ensure the root path has a trailing slash so that relative URLs work
|
// ensure the root path has a trailing slash so that relative URLs work
|
||||||
|
|
@ -806,7 +807,7 @@ func (s *Server) HandleGetList(w http.ResponseWriter, r *Request) {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
addr, err := s.api.Resolve(r.uri)
|
addr, err := s.api.Resolve(ctx, r.uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
getListFail.Inc(1)
|
getListFail.Inc(1)
|
||||||
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
||||||
|
|
@ -814,7 +815,7 @@ func (s *Server) HandleGetList(w http.ResponseWriter, r *Request) {
|
||||||
}
|
}
|
||||||
log.Debug("handle.get.list: resolved", "ruid", r.ruid, "key", addr)
|
log.Debug("handle.get.list: resolved", "ruid", r.ruid, "key", addr)
|
||||||
|
|
||||||
list, err := s.getManifestList(addr, r.uri.Path)
|
list, err := s.getManifestList(ctx, addr, r.uri.Path)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
getListFail.Inc(1)
|
getListFail.Inc(1)
|
||||||
|
|
@ -845,8 +846,8 @@ func (s *Server) HandleGetList(w http.ResponseWriter, r *Request) {
|
||||||
json.NewEncoder(w).Encode(&list)
|
json.NewEncoder(w).Encode(&list)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) getManifestList(addr storage.Address, prefix string) (list api.ManifestList, err error) {
|
func (s *Server) getManifestList(ctx context.Context, addr storage.Address, prefix string) (list api.ManifestList, err error) {
|
||||||
walker, err := s.api.NewManifestWalker(addr, nil)
|
walker, err := s.api.NewManifestWalker(ctx, addr, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -903,7 +904,7 @@ func (s *Server) getManifestList(addr storage.Address, prefix string) (list api.
|
||||||
|
|
||||||
// HandleGetFile handles a GET request to bzz://<manifest>/<path> and responds
|
// HandleGetFile handles a GET request to bzz://<manifest>/<path> and responds
|
||||||
// with the content of the file at <path> from the given <manifest>
|
// with the content of the file at <path> from the given <manifest>
|
||||||
func (s *Server) HandleGetFile(w http.ResponseWriter, r *Request) {
|
func (s *Server) HandleGetFile(ctx context.Context, w http.ResponseWriter, r *Request) {
|
||||||
log.Debug("handle.get.file", "ruid", r.ruid)
|
log.Debug("handle.get.file", "ruid", r.ruid)
|
||||||
getFileCount.Inc(1)
|
getFileCount.Inc(1)
|
||||||
// ensure the root path has a trailing slash so that relative URLs work
|
// ensure the root path has a trailing slash so that relative URLs work
|
||||||
|
|
@ -915,7 +916,7 @@ func (s *Server) HandleGetFile(w http.ResponseWriter, r *Request) {
|
||||||
manifestAddr := r.uri.Address()
|
manifestAddr := r.uri.Address()
|
||||||
|
|
||||||
if manifestAddr == nil {
|
if manifestAddr == nil {
|
||||||
manifestAddr, err = s.api.Resolve(r.uri)
|
manifestAddr, err = s.api.Resolve(ctx, r.uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
getFileFail.Inc(1)
|
getFileFail.Inc(1)
|
||||||
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
Respond(w, r, fmt.Sprintf("cannot resolve %s: %s", r.uri.Addr, err), http.StatusNotFound)
|
||||||
|
|
@ -927,7 +928,7 @@ func (s *Server) HandleGetFile(w http.ResponseWriter, r *Request) {
|
||||||
|
|
||||||
log.Debug("handle.get.file: resolved", "ruid", r.ruid, "key", manifestAddr)
|
log.Debug("handle.get.file: resolved", "ruid", r.ruid, "key", manifestAddr)
|
||||||
|
|
||||||
reader, contentType, status, contentKey, err := s.api.Get(manifestAddr, r.uri.Path)
|
reader, contentType, status, contentKey, err := s.api.Get(ctx, manifestAddr, r.uri.Path)
|
||||||
|
|
||||||
etag := common.Bytes2Hex(contentKey)
|
etag := common.Bytes2Hex(contentKey)
|
||||||
noneMatchEtag := r.Header.Get("If-None-Match")
|
noneMatchEtag := r.Header.Get("If-None-Match")
|
||||||
|
|
@ -954,7 +955,7 @@ func (s *Server) HandleGetFile(w http.ResponseWriter, r *Request) {
|
||||||
//the request results in ambiguous files
|
//the request results in ambiguous files
|
||||||
//e.g. /read with readme.md and readinglist.txt available in manifest
|
//e.g. /read with readme.md and readinglist.txt available in manifest
|
||||||
if status == http.StatusMultipleChoices {
|
if status == http.StatusMultipleChoices {
|
||||||
list, err := s.getManifestList(manifestAddr, r.uri.Path)
|
list, err := s.getManifestList(ctx, manifestAddr, r.uri.Path)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
getFileFail.Inc(1)
|
getFileFail.Inc(1)
|
||||||
|
|
@ -1011,6 +1012,8 @@ func (b bufferedReadSeeker) Seek(offset int64, whence int) (int64, error) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) ServeHTTP(rw http.ResponseWriter, r *http.Request) {
|
func (s *Server) ServeHTTP(rw http.ResponseWriter, r *http.Request) {
|
||||||
|
ctx := context.TODO()
|
||||||
|
|
||||||
defer metrics.GetOrRegisterResettingTimer(fmt.Sprintf("http.request.%s.time", r.Method), nil).UpdateSince(time.Now())
|
defer metrics.GetOrRegisterResettingTimer(fmt.Sprintf("http.request.%s.time", r.Method), nil).UpdateSince(time.Now())
|
||||||
req := &Request{Request: *r, ruid: uuid.New()[:8]}
|
req := &Request{Request: *r, ruid: uuid.New()[:8]}
|
||||||
metrics.GetOrRegisterCounter(fmt.Sprintf("http.request.%s", r.Method), nil).Inc(1)
|
metrics.GetOrRegisterCounter(fmt.Sprintf("http.request.%s", r.Method), nil).Inc(1)
|
||||||
|
|
@ -1055,16 +1058,16 @@ func (s *Server) ServeHTTP(rw http.ResponseWriter, r *http.Request) {
|
||||||
case "POST":
|
case "POST":
|
||||||
if uri.Raw() {
|
if uri.Raw() {
|
||||||
log.Debug("handlePostRaw")
|
log.Debug("handlePostRaw")
|
||||||
s.HandlePostRaw(w, req)
|
s.HandlePostRaw(ctx, w, req)
|
||||||
} else if uri.Resource() {
|
} else if uri.Resource() {
|
||||||
log.Debug("handlePostResource")
|
log.Debug("handlePostResource")
|
||||||
s.HandlePostResource(w, req)
|
s.HandlePostResource(ctx, w, req)
|
||||||
} else if uri.Immutable() || uri.List() || uri.Hash() {
|
} else if uri.Immutable() || uri.List() || uri.Hash() {
|
||||||
log.Debug("POST not allowed on immutable, list or hash")
|
log.Debug("POST not allowed on immutable, list or hash")
|
||||||
Respond(w, req, fmt.Sprintf("POST method on scheme %s not allowed", uri.Scheme), http.StatusMethodNotAllowed)
|
Respond(w, req, fmt.Sprintf("POST method on scheme %s not allowed", uri.Scheme), http.StatusMethodNotAllowed)
|
||||||
} else {
|
} else {
|
||||||
log.Debug("handlePostFiles")
|
log.Debug("handlePostFiles")
|
||||||
s.HandlePostFiles(w, req)
|
s.HandlePostFiles(ctx, w, req)
|
||||||
}
|
}
|
||||||
|
|
||||||
case "PUT":
|
case "PUT":
|
||||||
|
|
@ -1076,31 +1079,31 @@ func (s *Server) ServeHTTP(rw http.ResponseWriter, r *http.Request) {
|
||||||
Respond(w, req, fmt.Sprintf("DELETE method to %s not allowed", uri), http.StatusBadRequest)
|
Respond(w, req, fmt.Sprintf("DELETE method to %s not allowed", uri), http.StatusBadRequest)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
s.HandleDelete(w, req)
|
s.HandleDelete(ctx, w, req)
|
||||||
|
|
||||||
case "GET":
|
case "GET":
|
||||||
|
|
||||||
if uri.Resource() {
|
if uri.Resource() {
|
||||||
s.HandleGetResource(w, req)
|
s.HandleGetResource(ctx, w, req)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
if uri.Raw() || uri.Hash() {
|
if uri.Raw() || uri.Hash() {
|
||||||
s.HandleGet(w, req)
|
s.HandleGet(ctx, w, req)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
if uri.List() {
|
if uri.List() {
|
||||||
s.HandleGetList(w, req)
|
s.HandleGetList(ctx, w, req)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
if r.Header.Get("Accept") == "application/x-tar" {
|
if r.Header.Get("Accept") == "application/x-tar" {
|
||||||
s.HandleGetFiles(w, req)
|
s.HandleGetFiles(ctx, w, req)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
s.HandleGetFile(w, req)
|
s.HandleGetFile(ctx, w, req)
|
||||||
|
|
||||||
default:
|
default:
|
||||||
Respond(w, req, fmt.Sprintf("%s method is not supported", r.Method), http.StatusMethodNotAllowed)
|
Respond(w, req, fmt.Sprintf("%s method is not supported", r.Method), http.StatusMethodNotAllowed)
|
||||||
|
|
@ -1109,8 +1112,8 @@ func (s *Server) ServeHTTP(rw http.ResponseWriter, r *http.Request) {
|
||||||
log.Info("served response", "ruid", req.ruid, "code", w.statusCode)
|
log.Info("served response", "ruid", req.ruid, "code", w.statusCode)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) updateManifest(addr storage.Address, update func(mw *api.ManifestWriter) error) (storage.Address, error) {
|
func (s *Server) updateManifest(ctx context.Context, addr storage.Address, update func(mw *api.ManifestWriter) error) (storage.Address, error) {
|
||||||
mw, err := s.api.NewManifestWriter(addr, nil)
|
mw, err := s.api.NewManifestWriter(ctx, addr, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -18,6 +18,7 @@ package http
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"context"
|
||||||
"crypto/rand"
|
"crypto/rand"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
|
|
@ -383,7 +384,7 @@ func testBzzGetPath(encrypted bool, t *testing.T) {
|
||||||
for i, mf := range testmanifest {
|
for i, mf := range testmanifest {
|
||||||
reader[i] = bytes.NewReader([]byte(mf))
|
reader[i] = bytes.NewReader([]byte(mf))
|
||||||
var wait func()
|
var wait func()
|
||||||
addr[i], wait, err = srv.FileStore.Store(reader[i], int64(len(mf)), encrypted)
|
addr[i], wait, err = srv.FileStore.Store(context.TODO(), reader[i], int64(len(mf)), encrypted)
|
||||||
for j := i + 1; j < len(testmanifest); j++ {
|
for j := i + 1; j < len(testmanifest); j++ {
|
||||||
testmanifest[j] = strings.Replace(testmanifest[j], fmt.Sprintf("<key%v>", i), addr[i].Hex(), -1)
|
testmanifest[j] = strings.Replace(testmanifest[j], fmt.Sprintf("<key%v>", i), addr[i].Hex(), -1)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -18,6 +18,7 @@ package api
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
|
@ -61,20 +62,20 @@ type ManifestList struct {
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewManifest creates and stores a new, empty manifest
|
// NewManifest creates and stores a new, empty manifest
|
||||||
func (a *API) NewManifest(toEncrypt bool) (storage.Address, error) {
|
func (a *API) NewManifest(ctx context.Context, toEncrypt bool) (storage.Address, error) {
|
||||||
var manifest Manifest
|
var manifest Manifest
|
||||||
data, err := json.Marshal(&manifest)
|
data, err := json.Marshal(&manifest)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
key, wait, err := a.Store(bytes.NewReader(data), int64(len(data)), toEncrypt)
|
key, wait, err := a.Store(ctx, bytes.NewReader(data), int64(len(data)), toEncrypt)
|
||||||
wait()
|
wait()
|
||||||
return key, err
|
return key, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// Manifest hack for supporting Mutable Resource Updates from the bzz: scheme
|
// Manifest hack for supporting Mutable Resource Updates from the bzz: scheme
|
||||||
// see swarm/api/api.go:API.Get() for more information
|
// see swarm/api/api.go:API.Get() for more information
|
||||||
func (a *API) NewResourceManifest(resourceAddr string) (storage.Address, error) {
|
func (a *API) NewResourceManifest(ctx context.Context, resourceAddr string) (storage.Address, error) {
|
||||||
var manifest Manifest
|
var manifest Manifest
|
||||||
entry := ManifestEntry{
|
entry := ManifestEntry{
|
||||||
Hash: resourceAddr,
|
Hash: resourceAddr,
|
||||||
|
|
@ -85,7 +86,7 @@ func (a *API) NewResourceManifest(resourceAddr string) (storage.Address, error)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
key, _, err := a.Store(bytes.NewReader(data), int64(len(data)), false)
|
key, _, err := a.Store(ctx, bytes.NewReader(data), int64(len(data)), false)
|
||||||
return key, err
|
return key, err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -96,8 +97,8 @@ type ManifestWriter struct {
|
||||||
quitC chan bool
|
quitC chan bool
|
||||||
}
|
}
|
||||||
|
|
||||||
func (a *API) NewManifestWriter(addr storage.Address, quitC chan bool) (*ManifestWriter, error) {
|
func (a *API) NewManifestWriter(ctx context.Context, addr storage.Address, quitC chan bool) (*ManifestWriter, error) {
|
||||||
trie, err := loadManifest(a.fileStore, addr, quitC)
|
trie, err := loadManifest(ctx, a.fileStore, addr, quitC)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("error loading manifest %s: %s", addr, err)
|
return nil, fmt.Errorf("error loading manifest %s: %s", addr, err)
|
||||||
}
|
}
|
||||||
|
|
@ -105,9 +106,8 @@ func (a *API) NewManifestWriter(addr storage.Address, quitC chan bool) (*Manifes
|
||||||
}
|
}
|
||||||
|
|
||||||
// AddEntry stores the given data and adds the resulting key to the manifest
|
// AddEntry stores the given data and adds the resulting key to the manifest
|
||||||
func (m *ManifestWriter) AddEntry(data io.Reader, e *ManifestEntry) (storage.Address, error) {
|
func (m *ManifestWriter) AddEntry(ctx context.Context, data io.Reader, e *ManifestEntry) (storage.Address, error) {
|
||||||
|
key, _, err := m.api.Store(ctx, data, e.Size, m.trie.encrypted)
|
||||||
key, _, err := m.api.Store(data, e.Size, m.trie.encrypted)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
@ -136,8 +136,8 @@ type ManifestWalker struct {
|
||||||
quitC chan bool
|
quitC chan bool
|
||||||
}
|
}
|
||||||
|
|
||||||
func (a *API) NewManifestWalker(addr storage.Address, quitC chan bool) (*ManifestWalker, error) {
|
func (a *API) NewManifestWalker(ctx context.Context, addr storage.Address, quitC chan bool) (*ManifestWalker, error) {
|
||||||
trie, err := loadManifest(a.fileStore, addr, quitC)
|
trie, err := loadManifest(ctx, a.fileStore, addr, quitC)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("error loading manifest %s: %s", addr, err)
|
return nil, fmt.Errorf("error loading manifest %s: %s", addr, err)
|
||||||
}
|
}
|
||||||
|
|
@ -204,10 +204,10 @@ type manifestTrieEntry struct {
|
||||||
subtrie *manifestTrie
|
subtrie *manifestTrie
|
||||||
}
|
}
|
||||||
|
|
||||||
func loadManifest(fileStore *storage.FileStore, hash storage.Address, quitC chan bool) (trie *manifestTrie, err error) { // non-recursive, subtrees are downloaded on-demand
|
func loadManifest(ctx context.Context, fileStore *storage.FileStore, hash storage.Address, quitC chan bool) (trie *manifestTrie, err error) { // non-recursive, subtrees are downloaded on-demand
|
||||||
log.Trace("manifest lookup", "key", hash)
|
log.Trace("manifest lookup", "key", hash)
|
||||||
// retrieve manifest via FileStore
|
// retrieve manifest via FileStore
|
||||||
manifestReader, isEncrypted := fileStore.Retrieve(hash)
|
manifestReader, isEncrypted := fileStore.Retrieve(ctx, hash)
|
||||||
log.Trace("reader retrieved", "key", hash)
|
log.Trace("reader retrieved", "key", hash)
|
||||||
return readManifest(manifestReader, hash, fileStore, isEncrypted, quitC)
|
return readManifest(manifestReader, hash, fileStore, isEncrypted, quitC)
|
||||||
}
|
}
|
||||||
|
|
@ -382,7 +382,7 @@ func (mt *manifestTrie) recalcAndStore() error {
|
||||||
}
|
}
|
||||||
|
|
||||||
sr := bytes.NewReader(manifest)
|
sr := bytes.NewReader(manifest)
|
||||||
key, wait, err2 := mt.fileStore.Store(sr, int64(len(manifest)), mt.encrypted)
|
key, wait, err2 := mt.fileStore.Store(context.TODO(), sr, int64(len(manifest)), mt.encrypted)
|
||||||
wait()
|
wait()
|
||||||
mt.ref = key
|
mt.ref = key
|
||||||
return err2
|
return err2
|
||||||
|
|
@ -391,7 +391,7 @@ func (mt *manifestTrie) recalcAndStore() error {
|
||||||
func (mt *manifestTrie) loadSubTrie(entry *manifestTrieEntry, quitC chan bool) (err error) {
|
func (mt *manifestTrie) loadSubTrie(entry *manifestTrieEntry, quitC chan bool) (err error) {
|
||||||
if entry.subtrie == nil {
|
if entry.subtrie == nil {
|
||||||
hash := common.Hex2Bytes(entry.Hash)
|
hash := common.Hex2Bytes(entry.Hash)
|
||||||
entry.subtrie, err = loadManifest(mt.fileStore, hash, quitC)
|
entry.subtrie, err = loadManifest(context.TODO(), mt.fileStore, hash, quitC)
|
||||||
entry.Hash = "" // might not match, should be recalculated
|
entry.Hash = "" // might not match, should be recalculated
|
||||||
}
|
}
|
||||||
return
|
return
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,7 @@
|
||||||
package api
|
package api
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"path"
|
"path"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/swarm/storage"
|
"github.com/ethereum/go-ethereum/swarm/storage"
|
||||||
|
|
@ -45,8 +46,8 @@ func NewStorage(api *API) *Storage {
|
||||||
// its content type
|
// its content type
|
||||||
//
|
//
|
||||||
// DEPRECATED: Use the HTTP API instead
|
// DEPRECATED: Use the HTTP API instead
|
||||||
func (s *Storage) Put(content, contentType string, toEncrypt bool) (storage.Address, func(), error) {
|
func (s *Storage) Put(ctx context.Context, content string, contentType string, toEncrypt bool) (storage.Address, func(), error) {
|
||||||
return s.api.Put(content, contentType, toEncrypt)
|
return s.api.Put(ctx, content, contentType, toEncrypt)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get retrieves the content from bzzpath and reads the response in full
|
// Get retrieves the content from bzzpath and reads the response in full
|
||||||
|
|
@ -57,16 +58,16 @@ func (s *Storage) Put(content, contentType string, toEncrypt bool) (storage.Addr
|
||||||
// size is resp.Size
|
// size is resp.Size
|
||||||
//
|
//
|
||||||
// DEPRECATED: Use the HTTP API instead
|
// DEPRECATED: Use the HTTP API instead
|
||||||
func (s *Storage) Get(bzzpath string) (*Response, error) {
|
func (s *Storage) Get(ctx context.Context, bzzpath string) (*Response, error) {
|
||||||
uri, err := Parse(path.Join("bzz:/", bzzpath))
|
uri, err := Parse(path.Join("bzz:/", bzzpath))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
addr, err := s.api.Resolve(uri)
|
addr, err := s.api.Resolve(ctx, uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
reader, mimeType, status, _, err := s.api.Get(addr, uri.Path)
|
reader, mimeType, status, _, err := s.api.Get(ctx, addr, uri.Path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
@ -87,16 +88,16 @@ func (s *Storage) Get(bzzpath string) (*Response, error) {
|
||||||
// and merge on to it. creating an entry w conentType (mime)
|
// and merge on to it. creating an entry w conentType (mime)
|
||||||
//
|
//
|
||||||
// DEPRECATED: Use the HTTP API instead
|
// DEPRECATED: Use the HTTP API instead
|
||||||
func (s *Storage) Modify(rootHash, path, contentHash, contentType string) (newRootHash string, err error) {
|
func (s *Storage) Modify(ctx context.Context, rootHash, path, contentHash, contentType string) (newRootHash string, err error) {
|
||||||
uri, err := Parse("bzz:/" + rootHash)
|
uri, err := Parse("bzz:/" + rootHash)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
addr, err := s.api.Resolve(uri)
|
addr, err := s.api.Resolve(ctx, uri)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
addr, err = s.api.Modify(addr, path, contentHash, contentType)
|
addr, err = s.api.Modify(ctx, addr, path, contentHash, contentType)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,7 @@
|
||||||
package api
|
package api
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"testing"
|
"testing"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -31,7 +32,7 @@ func TestStoragePutGet(t *testing.T) {
|
||||||
content := "hello"
|
content := "hello"
|
||||||
exp := expResponse(content, "text/plain", 0)
|
exp := expResponse(content, "text/plain", 0)
|
||||||
// exp := expResponse([]byte(content), "text/plain", 0)
|
// exp := expResponse([]byte(content), "text/plain", 0)
|
||||||
bzzkey, wait, err := api.Put(content, exp.MimeType, toEncrypt)
|
bzzkey, wait, err := api.Put(context.TODO(), content, exp.MimeType, toEncrypt)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("unexpected error: %v", err)
|
t.Fatalf("unexpected error: %v", err)
|
||||||
}
|
}
|
||||||
|
|
@ -42,7 +43,7 @@ func TestStoragePutGet(t *testing.T) {
|
||||||
checkResponse(t, resp0, exp)
|
checkResponse(t, resp0, exp)
|
||||||
|
|
||||||
// check storage#Get
|
// check storage#Get
|
||||||
resp, err := api.Get(bzzhash)
|
resp, err := api.Get(context.TODO(), bzzhash)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("unexpected error: %v", err)
|
t.Fatalf("unexpected error: %v", err)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -84,7 +84,7 @@ func (sf *SwarmFile) Attr(ctx context.Context, a *fuse.Attr) error {
|
||||||
a.Gid = uint32(os.Getegid())
|
a.Gid = uint32(os.Getegid())
|
||||||
|
|
||||||
if sf.fileSize == -1 {
|
if sf.fileSize == -1 {
|
||||||
reader, _ := sf.mountInfo.swarmApi.Retrieve(sf.addr)
|
reader, _ := sf.mountInfo.swarmApi.Retrieve(ctx, sf.addr)
|
||||||
quitC := make(chan bool)
|
quitC := make(chan bool)
|
||||||
size, err := reader.Size(quitC)
|
size, err := reader.Size(quitC)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
@ -104,7 +104,7 @@ func (sf *SwarmFile) Read(ctx context.Context, req *fuse.ReadRequest, resp *fuse
|
||||||
sf.lock.RLock()
|
sf.lock.RLock()
|
||||||
defer sf.lock.RUnlock()
|
defer sf.lock.RUnlock()
|
||||||
if sf.reader == nil {
|
if sf.reader == nil {
|
||||||
sf.reader, _ = sf.mountInfo.swarmApi.Retrieve(sf.addr)
|
sf.reader, _ = sf.mountInfo.swarmApi.Retrieve(ctx, sf.addr)
|
||||||
}
|
}
|
||||||
buf := make([]byte, req.Size)
|
buf := make([]byte, req.Size)
|
||||||
n, err := sf.reader.ReadAt(buf, req.Offset)
|
n, err := sf.reader.ReadAt(buf, req.Offset)
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,7 @@ package fuse
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"context"
|
||||||
"crypto/rand"
|
"crypto/rand"
|
||||||
"flag"
|
"flag"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
|
@ -110,7 +111,7 @@ func createTestFilesAndUploadToSwarm(t *testing.T, api *api.API, files map[strin
|
||||||
}
|
}
|
||||||
|
|
||||||
//upload directory to swarm and return hash
|
//upload directory to swarm and return hash
|
||||||
bzzhash, err := api.Upload(uploadDir, "", toEncrypt)
|
bzzhash, err := api.Upload(context.TODO(), uploadDir, "", toEncrypt)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("Error uploading directory %v: %vm encryption: %v", uploadDir, err, toEncrypt)
|
t.Fatalf("Error uploading directory %v: %vm encryption: %v", uploadDir, err, toEncrypt)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -19,6 +19,7 @@
|
||||||
package fuse
|
package fuse
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"os"
|
"os"
|
||||||
|
|
@ -104,7 +105,7 @@ func (swarmfs *SwarmFS) Mount(mhash, mountpoint string) (*MountInfo, error) {
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Trace("swarmfs mount: getting manifest tree")
|
log.Trace("swarmfs mount: getting manifest tree")
|
||||||
_, manifestEntryMap, err := swarmfs.swarmApi.BuildDirectoryTree(mhash, true)
|
_, manifestEntryMap, err := swarmfs.swarmApi.BuildDirectoryTree(context.TODO(), mhash, true)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -47,7 +47,7 @@ func externalUnmount(mountPoint string) error {
|
||||||
}
|
}
|
||||||
|
|
||||||
func addFileToSwarm(sf *SwarmFile, content []byte, size int) error {
|
func addFileToSwarm(sf *SwarmFile, content []byte, size int) error {
|
||||||
fkey, mhash, err := sf.mountInfo.swarmApi.AddFile(sf.mountInfo.LatestManifest, sf.path, sf.name, content, true)
|
fkey, mhash, err := sf.mountInfo.swarmApi.AddFile(context.TODO(), sf.mountInfo.LatestManifest, sf.path, sf.name, content, true)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
@ -66,7 +66,7 @@ func addFileToSwarm(sf *SwarmFile, content []byte, size int) error {
|
||||||
}
|
}
|
||||||
|
|
||||||
func removeFileFromSwarm(sf *SwarmFile) error {
|
func removeFileFromSwarm(sf *SwarmFile) error {
|
||||||
mkey, err := sf.mountInfo.swarmApi.RemoveFile(sf.mountInfo.LatestManifest, sf.path, sf.name, true)
|
mkey, err := sf.mountInfo.swarmApi.RemoveFile(context.TODO(), sf.mountInfo.LatestManifest, sf.path, sf.name, true)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
@ -102,7 +102,7 @@ func removeDirectoryFromSwarm(sd *SwarmDir) error {
|
||||||
}
|
}
|
||||||
|
|
||||||
func appendToExistingFileInSwarm(sf *SwarmFile, content []byte, offset int64, length int64) error {
|
func appendToExistingFileInSwarm(sf *SwarmFile, content []byte, offset int64, length int64) error {
|
||||||
fkey, mhash, err := sf.mountInfo.swarmApi.AppendFile(sf.mountInfo.LatestManifest, sf.path, sf.name, sf.fileSize, content, sf.addr, offset, length, true)
|
fkey, mhash, err := sf.mountInfo.swarmApi.AppendFile(context.TODO(), sf.mountInfo.LatestManifest, sf.path, sf.name, sf.fileSize, content, sf.addr, offset, length, true)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -250,7 +250,7 @@ func (r *TestRegistry) APIs() []rpc.API {
|
||||||
}
|
}
|
||||||
|
|
||||||
func readAll(fileStore *storage.FileStore, hash []byte) (int64, error) {
|
func readAll(fileStore *storage.FileStore, hash []byte) (int64, error) {
|
||||||
r, _ := fileStore.Retrieve(hash)
|
r, _ := fileStore.Retrieve(context.TODO(), hash)
|
||||||
buf := make([]byte, 1024)
|
buf := make([]byte, 1024)
|
||||||
var n int
|
var n int
|
||||||
var total int64
|
var total int64
|
||||||
|
|
|
||||||
|
|
@ -345,7 +345,7 @@ func testDeliveryFromNodes(t *testing.T, nodes, conns, chunkCount int, skipCheck
|
||||||
// here we distribute chunks of a random file into Stores of nodes 1 to nodes
|
// here we distribute chunks of a random file into Stores of nodes 1 to nodes
|
||||||
rrFileStore := storage.NewFileStore(newRoundRobinStore(sim.Stores[1:]...), storage.NewFileStoreParams())
|
rrFileStore := storage.NewFileStore(newRoundRobinStore(sim.Stores[1:]...), storage.NewFileStoreParams())
|
||||||
size := chunkCount * chunkSize
|
size := chunkCount * chunkSize
|
||||||
fileHash, wait, err := rrFileStore.Store(io.LimitReader(crand.Reader, int64(size)), int64(size), false)
|
fileHash, wait, err := rrFileStore.Store(context.TODO(), io.LimitReader(crand.Reader, int64(size)), int64(size), false)
|
||||||
// wait until all chunks stored
|
// wait until all chunks stored
|
||||||
wait()
|
wait()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
@ -627,7 +627,7 @@ Loop:
|
||||||
hashes := make([]storage.Address, chunkCount)
|
hashes := make([]storage.Address, chunkCount)
|
||||||
for i := 0; i < chunkCount; i++ {
|
for i := 0; i < chunkCount; i++ {
|
||||||
// create actual size real chunks
|
// create actual size real chunks
|
||||||
hash, wait, err := remoteFileStore.Store(io.LimitReader(crand.Reader, int64(chunkSize)), int64(chunkSize), false)
|
hash, wait, err := remoteFileStore.Store(context.TODO(), io.LimitReader(crand.Reader, int64(chunkSize)), int64(chunkSize), false)
|
||||||
// wait until all chunks stored
|
// wait until all chunks stored
|
||||||
wait()
|
wait()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
|
||||||
|
|
@ -117,7 +117,7 @@ func testIntervals(t *testing.T, live bool, history *Range, skipCheck bool) {
|
||||||
|
|
||||||
fileStore := storage.NewFileStore(sim.Stores[0], storage.NewFileStoreParams())
|
fileStore := storage.NewFileStore(sim.Stores[0], storage.NewFileStoreParams())
|
||||||
size := chunkCount * chunkSize
|
size := chunkCount * chunkSize
|
||||||
_, wait, err := fileStore.Store(io.LimitReader(crand.Reader, int64(size)), int64(size), false)
|
_, wait, err := fileStore.Store(context.TODO(), io.LimitReader(crand.Reader, int64(size)), int64(size), false)
|
||||||
wait()
|
wait()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
|
|
|
||||||
|
|
@ -410,7 +410,7 @@ func runFileRetrievalTest(nodeCount int) error {
|
||||||
fileStore := registries[id].fileStore
|
fileStore := registries[id].fileStore
|
||||||
//check all chunks
|
//check all chunks
|
||||||
for i, hash := range conf.hashes {
|
for i, hash := range conf.hashes {
|
||||||
reader, _ := fileStore.Retrieve(hash)
|
reader, _ := fileStore.Retrieve(context.TODO(), hash)
|
||||||
//check that we can read the file size and that it corresponds to the generated file size
|
//check that we can read the file size and that it corresponds to the generated file size
|
||||||
if s, err := reader.Size(nil); err != nil || s != int64(len(randomFiles[i])) {
|
if s, err := reader.Size(nil); err != nil || s != int64(len(randomFiles[i])) {
|
||||||
allSuccess = false
|
allSuccess = false
|
||||||
|
|
@ -697,7 +697,7 @@ func runRetrievalTest(chunkCount int, nodeCount int) error {
|
||||||
fileStore := registries[id].fileStore
|
fileStore := registries[id].fileStore
|
||||||
//check all chunks
|
//check all chunks
|
||||||
for _, chnk := range conf.hashes {
|
for _, chnk := range conf.hashes {
|
||||||
reader, _ := fileStore.Retrieve(chnk)
|
reader, _ := fileStore.Retrieve(context.TODO(), chnk)
|
||||||
//assuming that reading the Size of the chunk is enough to know we found it
|
//assuming that reading the Size of the chunk is enough to know we found it
|
||||||
if s, err := reader.Size(nil); err != nil || s != chunkSize {
|
if s, err := reader.Size(nil); err != nil || s != chunkSize {
|
||||||
allSuccess = false
|
allSuccess = false
|
||||||
|
|
@ -765,7 +765,7 @@ func uploadFilesToNodes(nodes []*simulations.Node) ([]storage.Address, []string,
|
||||||
return nil, nil, err
|
return nil, nil, err
|
||||||
}
|
}
|
||||||
//store it (upload it) on the FileStore
|
//store it (upload it) on the FileStore
|
||||||
rk, wait, err := fileStore.Store(strings.NewReader(rfiles[i]), int64(len(rfiles[i])), false)
|
rk, wait, err := fileStore.Store(context.TODO(), strings.NewReader(rfiles[i]), int64(len(rfiles[i])), false)
|
||||||
log.Debug("Uploaded random string file to node")
|
log.Debug("Uploaded random string file to node")
|
||||||
wait()
|
wait()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
|
||||||
|
|
@ -581,7 +581,7 @@ func uploadFileToSingleNodeStore(id discover.NodeID, chunkCount int) ([]storage.
|
||||||
fileStore := storage.NewFileStore(lstore, storage.NewFileStoreParams())
|
fileStore := storage.NewFileStore(lstore, storage.NewFileStoreParams())
|
||||||
var rootAddrs []storage.Address
|
var rootAddrs []storage.Address
|
||||||
for i := 0; i < chunkCount; i++ {
|
for i := 0; i < chunkCount; i++ {
|
||||||
rk, wait, err := fileStore.Store(io.LimitReader(crand.Reader, int64(size)), int64(size), false)
|
rk, wait, err := fileStore.Store(context.TODO(), io.LimitReader(crand.Reader, int64(size)), int64(size), false)
|
||||||
wait()
|
wait()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
|
|
|
||||||
|
|
@ -202,7 +202,7 @@ func testSyncBetweenNodes(t *testing.T, nodes, conns, chunkCount int, skipCheck
|
||||||
// here we distribute chunks of a random file into stores 1...nodes
|
// here we distribute chunks of a random file into stores 1...nodes
|
||||||
rrFileStore := storage.NewFileStore(newRoundRobinStore(sim.Stores[1:]...), storage.NewFileStoreParams())
|
rrFileStore := storage.NewFileStore(newRoundRobinStore(sim.Stores[1:]...), storage.NewFileStoreParams())
|
||||||
size := chunkCount * chunkSize
|
size := chunkCount * chunkSize
|
||||||
_, wait, err := rrFileStore.Store(io.LimitReader(crand.Reader, int64(size)), int64(size), false)
|
_, wait, err := rrFileStore.Store(context.TODO(), io.LimitReader(crand.Reader, int64(size)), int64(size), false)
|
||||||
// need to wait cos we then immediately collect the relevant bin content
|
// need to wait cos we then immediately collect the relevant bin content
|
||||||
wait()
|
wait()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
|
||||||
|
|
@ -508,7 +508,7 @@ func uploadFile(swarm *Swarm) (storage.Address, string, error) {
|
||||||
// File data is very short, but it is ensured that its
|
// File data is very short, but it is ensured that its
|
||||||
// uniqueness is very certain.
|
// uniqueness is very certain.
|
||||||
data := fmt.Sprintf("test content %s %x", time.Now().Round(0), b)
|
data := fmt.Sprintf("test content %s %x", time.Now().Round(0), b)
|
||||||
k, wait, err := swarm.api.Put(data, "text/plain", false)
|
k, wait, err := swarm.api.Put(context.TODO(), data, "text/plain", false)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, "", err
|
return nil, "", err
|
||||||
}
|
}
|
||||||
|
|
@ -570,7 +570,7 @@ func retrieve(
|
||||||
|
|
||||||
log.Debug("api get: check file", "node", id.String(), "key", f.addr.String(), "total files found", atomic.LoadUint64(totalFoundCount))
|
log.Debug("api get: check file", "node", id.String(), "key", f.addr.String(), "total files found", atomic.LoadUint64(totalFoundCount))
|
||||||
|
|
||||||
r, _, _, _, err := swarm.api.Get(f.addr, "/")
|
r, _, _, _, err := swarm.api.Get(context.TODO(), f.addr, "/")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
errc <- fmt.Errorf("api get: node %s, key %s, kademlia %s: %v", id, f.addr, swarm.bzz.Hive, err)
|
errc <- fmt.Errorf("api get: node %s, key %s, kademlia %s: %v", id, f.addr, swarm.bzz.Hive, err)
|
||||||
return
|
return
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,7 @@
|
||||||
package storage
|
package storage
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"io"
|
"io"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -78,7 +79,7 @@ func NewFileStore(store ChunkStore, params *FileStoreParams) *FileStore {
|
||||||
// Chunk retrieval blocks on netStore requests with a timeout so reader will
|
// Chunk retrieval blocks on netStore requests with a timeout so reader will
|
||||||
// report error if retrieval of chunks within requested range time out.
|
// report error if retrieval of chunks within requested range time out.
|
||||||
// It returns a reader with the chunk data and whether the content was encrypted
|
// It returns a reader with the chunk data and whether the content was encrypted
|
||||||
func (f *FileStore) Retrieve(addr Address) (reader *LazyChunkReader, isEncrypted bool) {
|
func (f *FileStore) Retrieve(ctx context.Context, addr Address) (reader *LazyChunkReader, isEncrypted bool) {
|
||||||
isEncrypted = len(addr) > f.hashFunc().Size()
|
isEncrypted = len(addr) > f.hashFunc().Size()
|
||||||
getter := NewHasherStore(f.ChunkStore, f.hashFunc, isEncrypted)
|
getter := NewHasherStore(f.ChunkStore, f.hashFunc, isEncrypted)
|
||||||
reader = TreeJoin(addr, getter, 0)
|
reader = TreeJoin(addr, getter, 0)
|
||||||
|
|
@ -87,7 +88,7 @@ func (f *FileStore) Retrieve(addr Address) (reader *LazyChunkReader, isEncrypted
|
||||||
|
|
||||||
// Public API. Main entry point for document storage directly. Used by the
|
// Public API. Main entry point for document storage directly. Used by the
|
||||||
// FS-aware API and httpaccess
|
// FS-aware API and httpaccess
|
||||||
func (f *FileStore) Store(data io.Reader, size int64, toEncrypt bool) (addr Address, wait func(), err error) {
|
func (f *FileStore) Store(ctx context.Context, data io.Reader, size int64, toEncrypt bool) (addr Address, wait func(), err error) {
|
||||||
putter := NewHasherStore(f.ChunkStore, f.hashFunc, toEncrypt)
|
putter := NewHasherStore(f.ChunkStore, f.hashFunc, toEncrypt)
|
||||||
return PyramidSplit(data, putter, putter)
|
return PyramidSplit(data, putter, putter)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -18,6 +18,7 @@ package storage
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"context"
|
||||||
"io"
|
"io"
|
||||||
"io/ioutil"
|
"io/ioutil"
|
||||||
"os"
|
"os"
|
||||||
|
|
@ -49,12 +50,12 @@ func testFileStoreRandom(toEncrypt bool, t *testing.T) {
|
||||||
defer os.RemoveAll("/tmp/bzz")
|
defer os.RemoveAll("/tmp/bzz")
|
||||||
|
|
||||||
reader, slice := generateRandomData(testDataSize)
|
reader, slice := generateRandomData(testDataSize)
|
||||||
key, wait, err := fileStore.Store(reader, testDataSize, toEncrypt)
|
key, wait, err := fileStore.Store(context.TODO(), reader, testDataSize, toEncrypt)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Errorf("Store error: %v", err)
|
t.Errorf("Store error: %v", err)
|
||||||
}
|
}
|
||||||
wait()
|
wait()
|
||||||
resultReader, isEncrypted := fileStore.Retrieve(key)
|
resultReader, isEncrypted := fileStore.Retrieve(context.TODO(), key)
|
||||||
if isEncrypted != toEncrypt {
|
if isEncrypted != toEncrypt {
|
||||||
t.Fatalf("isEncrypted expected %v got %v", toEncrypt, isEncrypted)
|
t.Fatalf("isEncrypted expected %v got %v", toEncrypt, isEncrypted)
|
||||||
}
|
}
|
||||||
|
|
@ -72,7 +73,7 @@ func testFileStoreRandom(toEncrypt bool, t *testing.T) {
|
||||||
ioutil.WriteFile("/tmp/slice.bzz.16M", slice, 0666)
|
ioutil.WriteFile("/tmp/slice.bzz.16M", slice, 0666)
|
||||||
ioutil.WriteFile("/tmp/result.bzz.16M", resultSlice, 0666)
|
ioutil.WriteFile("/tmp/result.bzz.16M", resultSlice, 0666)
|
||||||
localStore.memStore = NewMemStore(NewDefaultStoreParams(), db)
|
localStore.memStore = NewMemStore(NewDefaultStoreParams(), db)
|
||||||
resultReader, isEncrypted = fileStore.Retrieve(key)
|
resultReader, isEncrypted = fileStore.Retrieve(context.TODO(), key)
|
||||||
if isEncrypted != toEncrypt {
|
if isEncrypted != toEncrypt {
|
||||||
t.Fatalf("isEncrypted expected %v got %v", toEncrypt, isEncrypted)
|
t.Fatalf("isEncrypted expected %v got %v", toEncrypt, isEncrypted)
|
||||||
}
|
}
|
||||||
|
|
@ -110,12 +111,12 @@ func testFileStoreCapacity(toEncrypt bool, t *testing.T) {
|
||||||
}
|
}
|
||||||
fileStore := NewFileStore(localStore, NewFileStoreParams())
|
fileStore := NewFileStore(localStore, NewFileStoreParams())
|
||||||
reader, slice := generateRandomData(testDataSize)
|
reader, slice := generateRandomData(testDataSize)
|
||||||
key, wait, err := fileStore.Store(reader, testDataSize, toEncrypt)
|
key, wait, err := fileStore.Store(context.TODO(), reader, testDataSize, toEncrypt)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Errorf("Store error: %v", err)
|
t.Errorf("Store error: %v", err)
|
||||||
}
|
}
|
||||||
wait()
|
wait()
|
||||||
resultReader, isEncrypted := fileStore.Retrieve(key)
|
resultReader, isEncrypted := fileStore.Retrieve(context.TODO(), key)
|
||||||
if isEncrypted != toEncrypt {
|
if isEncrypted != toEncrypt {
|
||||||
t.Fatalf("isEncrypted expected %v got %v", toEncrypt, isEncrypted)
|
t.Fatalf("isEncrypted expected %v got %v", toEncrypt, isEncrypted)
|
||||||
}
|
}
|
||||||
|
|
@ -134,7 +135,7 @@ func testFileStoreCapacity(toEncrypt bool, t *testing.T) {
|
||||||
memStore.setCapacity(0)
|
memStore.setCapacity(0)
|
||||||
// check whether it is, indeed, empty
|
// check whether it is, indeed, empty
|
||||||
fileStore.ChunkStore = memStore
|
fileStore.ChunkStore = memStore
|
||||||
resultReader, isEncrypted = fileStore.Retrieve(key)
|
resultReader, isEncrypted = fileStore.Retrieve(context.TODO(), key)
|
||||||
if isEncrypted != toEncrypt {
|
if isEncrypted != toEncrypt {
|
||||||
t.Fatalf("isEncrypted expected %v got %v", toEncrypt, isEncrypted)
|
t.Fatalf("isEncrypted expected %v got %v", toEncrypt, isEncrypted)
|
||||||
}
|
}
|
||||||
|
|
@ -144,7 +145,7 @@ func testFileStoreCapacity(toEncrypt bool, t *testing.T) {
|
||||||
// check how it works with localStore
|
// check how it works with localStore
|
||||||
fileStore.ChunkStore = localStore
|
fileStore.ChunkStore = localStore
|
||||||
// localStore.dbStore.setCapacity(0)
|
// localStore.dbStore.setCapacity(0)
|
||||||
resultReader, isEncrypted = fileStore.Retrieve(key)
|
resultReader, isEncrypted = fileStore.Retrieve(context.TODO(), key)
|
||||||
if isEncrypted != toEncrypt {
|
if isEncrypted != toEncrypt {
|
||||||
t.Fatalf("isEncrypted expected %v got %v", toEncrypt, isEncrypted)
|
t.Fatalf("isEncrypted expected %v got %v", toEncrypt, isEncrypted)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,7 @@
|
||||||
package swarm
|
package swarm
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"encoding/hex"
|
"encoding/hex"
|
||||||
"io/ioutil"
|
"io/ioutil"
|
||||||
"math/rand"
|
"math/rand"
|
||||||
|
|
@ -347,7 +348,7 @@ func testLocalStoreAndRetrieve(t *testing.T, swarm *Swarm, n int, randomData boo
|
||||||
}
|
}
|
||||||
dataPut := string(slice)
|
dataPut := string(slice)
|
||||||
|
|
||||||
k, wait, err := swarm.api.Store(strings.NewReader(dataPut), int64(len(dataPut)), false)
|
k, wait, err := swarm.api.Store(context.TODO(), strings.NewReader(dataPut), int64(len(dataPut)), false)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
|
|
@ -355,7 +356,7 @@ func testLocalStoreAndRetrieve(t *testing.T, swarm *Swarm, n int, randomData boo
|
||||||
wait()
|
wait()
|
||||||
}
|
}
|
||||||
|
|
||||||
r, _ := swarm.api.Retrieve(k)
|
r, _ := swarm.api.Retrieve(context.TODO(), k)
|
||||||
|
|
||||||
d, err := ioutil.ReadAll(r)
|
d, err := ioutil.ReadAll(r)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue