This commit is contained in:
Todd Gruben 2016-02-26 12:13:03 -06:00
commit 7c33b7fcf2
4 changed files with 215 additions and 0 deletions

View file

@ -26,6 +26,9 @@ const (
// SnapshotExt is the file extension used for an in-process snapshot.
SnapshotExt = ".snapshotting"
// CopyExt is the file extension used for the temp file used while copying.
CopyExt = ".copying"
// CacheExt is the file extension for persisted cache ids.
CacheExt = ".cache"
@ -644,6 +647,74 @@ func (f *Fragment) flushCache() error {
return nil
}
// WriteTo writes the fragment's data to w.
func (f *Fragment) WriteTo(w io.Writer) (n int64, err error) {
// Open separate file descriptor to read from.
file, err := os.Open(f.path)
if err != nil {
return 0, err
}
defer file.Close()
// Retrieve the current file size under lock so we don't read
// while an operation is appending to the end.
var sz int64
if err := func() error {
f.mu.Lock()
defer f.mu.Unlock()
fi, err := file.Stat()
if err != nil {
return err
}
sz = fi.Size()
return nil
}(); err != nil {
return 0, err
}
// Copy the file up to the last known size.
// This is done outside the lock because the storage format is append-only.
return io.CopyN(w, file, sz)
}
// ReadFrom reads a data file from r and loads it into the fragment.
func (f *Fragment) ReadFrom(r io.Reader) (n int64, err error) {
f.mu.Lock()
defer f.mu.Unlock()
// Create a temporary file to copy into.
path := f.path + CopyExt
file, err := os.Create(path)
if err != nil {
return 0, err
}
defer file.Close()
// Copy reader into temporary path.
if n, err = io.Copy(file, r); err != nil {
return n, err
}
// Close current storage.
if err := f.closeStorage(); err != nil {
return n, err
}
// Move snapshot to data file location.
if err := os.Rename(path, f.path); err != nil {
return n, err
}
// Reopen storage.
if err := f.openStorage(); err != nil {
return n, err
}
return n, nil
}
func madvise(b []byte, advice int) (err error) {
_, _, e1 := syscall.Syscall(syscall.SYS_MADVISE, uintptr(unsafe.Pointer(&b[0])), uintptr(len(b)), uintptr(advice))
if e1 != 0 {

View file

@ -1,6 +1,7 @@
package pilosa_test
import (
"bytes"
"io/ioutil"
"os"
"reflect"
@ -279,6 +280,48 @@ func TestFragment_RankCache_Persistence(t *testing.T) {
}
}
// Ensure a fragment can be copied to another fragment.
func TestFragment_WriteTo_ReadFrom(t *testing.T) {
f0 := MustOpenFragment("d", "f", 0)
defer f0.Close()
// Set and then clear bits on the fragment.
if _, err := f0.SetBit(1000, 1, nil, 0); err != nil {
t.Fatal(err)
} else if _, err := f0.SetBit(1000, 2, nil, 0); err != nil {
t.Fatal(err)
} else if _, err := f0.ClearBit(1000, 1); err != nil {
t.Fatal(err)
}
// Write fragment to a buffer.
var buf bytes.Buffer
wn, err := f0.WriteTo(&buf)
if err != nil {
t.Fatal(err)
}
// Read into another fragment.
f1 := MustOpenFragment("d", "f", 0)
if rn, err := f1.ReadFrom(&buf); err != nil {
t.Fatal(err)
} else if wn != rn {
t.Fatalf("read/write byte count mismatch: wn=%d, rn=%d", wn, rn)
}
// Verify data in other fragment.
if a := f1.Bitmap(1000).Bits(); !reflect.DeepEqual(a, []uint64{2}) {
t.Fatalf("unexpected bits: %+v", a)
}
// Close and reopen the fragment & verify the data.
if err := f1.Reopen(); err != nil {
t.Fatal(err)
} else if a := f1.Bitmap(1000).Bits(); !reflect.DeepEqual(a, []uint64{2}) {
t.Fatalf("unexpected bits: %+v", a)
}
}
// Fragment is a test wrapper for pilosa.Fragment.
type Fragment struct {
*pilosa.Fragment

View file

@ -80,6 +80,15 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
default:
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
}
case "/fragment/data":
switch r.Method {
case "GET":
h.handleGetFragmentData(w, r)
case "POST":
h.handlePostFragmentData(w, r)
default:
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
}
case "/version":
h.handleVersion(w, r)
@ -353,6 +362,7 @@ func (h *Handler) handleGetSlicesNodes(w http.ResponseWriter, r *http.Request) {
slice, err := strconv.ParseUint(q.Get("slice"), 10, 64)
if err != nil {
http.Error(w, "slice required", http.StatusBadRequest)
return
}
// Retrieve slice owner nodes.
@ -364,6 +374,53 @@ func (h *Handler) handleGetSlicesNodes(w http.ResponseWriter, r *http.Request) {
}
}
// handleGetFragmentBackup handles GET /fragment/data requests.
func (h *Handler) handleGetFragmentData(w http.ResponseWriter, r *http.Request) {
// Read slice parameter.
q := r.URL.Query()
slice, err := strconv.ParseUint(q.Get("slice"), 10, 64)
if err != nil {
http.Error(w, "slice required", http.StatusBadRequest)
return
}
// Retrieve fragment from index.
f := h.Index.Fragment(q.Get("db"), q.Get("frame"), slice)
if f == nil {
http.Error(w, "fragment not found", http.StatusNotFound)
return
}
// Stream fragment to response body.
if _, err := f.WriteTo(w); err != nil {
h.logger().Printf("fragment backup error: %s", err)
}
}
// handlePostFragmentRestore handles POST /fragment/data requests.
func (h *Handler) handlePostFragmentData(w http.ResponseWriter, r *http.Request) {
// Read slice parameter.
q := r.URL.Query()
slice, err := strconv.ParseUint(q.Get("slice"), 10, 64)
if err != nil {
http.Error(w, "slice required", http.StatusBadRequest)
return
}
// Retrieve fragment from index.
f, err := h.Index.CreateFragmentIfNotExists(q.Get("db"), q.Get("frame"), slice)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Read fragment in from request body.
if _, err := f.ReadFrom(r.Body); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
// handleGetVersion handles /version requests.
func (h *Handler) handleVersion(w http.ResponseWriter, r *http.Request) {
if err := json.NewEncoder(w).Encode(struct {

View file

@ -389,6 +389,50 @@ func TestHandler_Query_ErrParse(t *testing.T) {
}
}
// Ensure the handler can backup a fragment and then restore it.
func TestHandler_Fragment_BackupRestore(t *testing.T) {
idx := MustOpenIndex()
defer idx.Close()
s := NewServer()
s.Handler.Index = idx.Index
defer s.Close()
// Set bits in the index.
f0 := idx.MustCreateFragmentIfNotExists("d", "f", 0)
f0.MustSetBits(100, 1, 2, 3)
// Begin backing up from slice d/f/0.
resp, err := http.Get(s.URL + "/fragment/data?db=d&frame=f&slice=0")
if err != nil {
t.Fatal(err)
}
defer resp.Body.Close()
// Ensure response came back OK.
if resp.StatusCode != http.StatusOK {
t.Fatalf("unexpected backup status code: %d", resp.StatusCode)
}
// Restore backup to slice x/y/0.
if resp, err := http.Post(s.URL+"/fragment/data?db=x&frame=y&slice=0", "application/octet-stream", resp.Body); err != nil {
t.Fatal(err)
} else if resp.StatusCode != http.StatusOK {
resp.Body.Close()
t.Fatalf("unexpected restore status code: %d", resp.StatusCode)
} else {
resp.Body.Close()
}
// Verify data is correctly restored.
f1 := idx.Fragment("x", "y", 0)
if f1 == nil {
t.Fatal("fragment x/y/0 not created")
} else if bits := f1.Bitmap(100).Bits(); !reflect.DeepEqual(bits, []uint64{1, 2, 3}) {
t.Fatalf("unexpected restored bits: %+v", bits)
}
}
// Ensure the handler can retrieve the version.
func TestHandler_Version(t *testing.T) {
h := NewHandler()