From c64199a69e7d38a8469e664e16fa02d6f7832eee Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Thu, 25 Feb 2016 17:29:45 -0700 Subject: [PATCH] add backup/restore endpoints This commit adds io.ReaderFrom and io.WriterTo implementations to the Fragment and also adds backup & restore endpoints to the Handler. --- fragment.go | 71 ++++++++++++++++++++++++++++++++++++++++++++++++ fragment_test.go | 43 +++++++++++++++++++++++++++++ handler.go | 57 ++++++++++++++++++++++++++++++++++++++ handler_test.go | 44 ++++++++++++++++++++++++++++++ 4 files changed, 215 insertions(+) diff --git a/fragment.go b/fragment.go index 6c5b1b7f6..c0f1bef0d 100644 --- a/fragment.go +++ b/fragment.go @@ -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" @@ -635,6 +638,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 { diff --git a/fragment_test.go b/fragment_test.go index e5903c124..8dc13ef05 100644 --- a/fragment_test.go +++ b/fragment_test.go @@ -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 diff --git a/handler.go b/handler.go index 52940e7e2..9c1896202 100644 --- a/handler.go +++ b/handler.go @@ -73,6 +73,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) @@ -328,6 +337,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. @@ -339,6 +349,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 { diff --git a/handler_test.go b/handler_test.go index 65d000c3e..c07f3e597 100644 --- a/handler_test.go +++ b/handler_test.go @@ -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()