diff --git a/main/main_async.go b/examples/async/async.go similarity index 98% rename from main/main_async.go rename to examples/async/async.go index 27c8eb8..d1562d2 100644 --- a/main/main_async.go +++ b/examples/async/async.go @@ -8,7 +8,7 @@ import ( "github.com/jhunters/bigqueue" ) -func main2() { +func main() { var queue = new(bigqueue.FileQueue) var DefaultOptions = &bigqueue.Options{ diff --git a/examples/async/go.mod b/examples/async/go.mod new file mode 100644 index 0000000..382de24 --- /dev/null +++ b/examples/async/go.mod @@ -0,0 +1,3 @@ +module async + +go 1.13 diff --git a/examples/nomal/go.mod b/examples/nomal/go.mod new file mode 100644 index 0000000..3e6a1ea --- /dev/null +++ b/examples/nomal/go.mod @@ -0,0 +1,5 @@ +module normal + +go 1.13 + +require github.com/jhunters/bigqueue v1.0.1 // indirect diff --git a/examples/nomal/go.sum b/examples/nomal/go.sum new file mode 100644 index 0000000..6b5c505 --- /dev/null +++ b/examples/nomal/go.sum @@ -0,0 +1,2 @@ +github.com/jhunters/bigqueue v1.0.1 h1:xIcBuzfm1ILMdGY5M+fukI84sqv87SPFUDEFT7xaKmc= +github.com/jhunters/bigqueue v1.0.1/go.mod h1:CbEObWKPe9f5OYtOu99fb92F+gJJewLhxeDm/V+WBio= diff --git a/examples/nomal/normal.go b/examples/nomal/normal.go new file mode 100644 index 0000000..dcd2eff --- /dev/null +++ b/examples/nomal/normal.go @@ -0,0 +1,55 @@ +package main + +import ( + "fmt" + "strconv" + + "github.com/jhunters/bigqueue" +) + +// a demo to show how to enqueue and dequeue data +func main() { + + var queue = new(bigqueue.FileQueue) + + // use custom options + var DefaultOptions = &bigqueue.Options{ + DataPageSize: bigqueue.DefaultDataPageSize, + GcLock: false, + IndexItemsPerPage: bigqueue.DefaultIndexItemsPerPage, + } + + // open queue files + err := queue.Open("./bin", "testqueue", DefaultOptions) + + if err != nil { + fmt.Println(err) + } + defer queue.Close() + + // do enqueue + for i := 1; i < 10; i++ { + data := []byte("hello jhunters" + strconv.Itoa(i)) + i, err := queue.Enqueue(data) + if err != nil { + fmt.Println(err) + } else { + fmt.Println("Enqueued index=", i, string(data)) + } + } + // do dequeue + for i := 1; i < 10; i++ { + index, bb, err := queue.Dequeue() + if err != nil { + fmt.Println(err) + } + if index != -1 { + fmt.Println(index, string(bb)) + } + + } + + // do gc action to free old data + queue.Gc() + +} diff --git a/main/main.go b/examples/subscribe/subscribe.go similarity index 57% rename from main/main.go rename to examples/subscribe/subscribe.go index 85158ef..0169829 100644 --- a/main/main.go +++ b/examples/subscribe/subscribe.go @@ -8,53 +8,6 @@ import ( "github.com/jhunters/bigqueue" ) -// a demo to show how to enqueue and dequeue data -func main() { - - var queue = new(bigqueue.FileQueue) - - // use custom options - var DefaultOptions = &bigqueue.Options{ - DataPageSize: bigqueue.DefaultDataPageSize, - GcLock: false, - IndexItemsPerPage: bigqueue.DefaultIndexItemsPerPage, - } - - // open queue files - err := queue.Open("./bin", "testqueue", DefaultOptions) - - if err != nil { - fmt.Println(err) - } - defer queue.Close() - - // do enqueue - for i := 1; i < 10; i++ { - data := []byte("hello jhunters" + strconv.Itoa(i)) - i, err := queue.Enqueue(data) - if err != nil { - fmt.Println(err) - } else { - fmt.Println("Enqueued index=", i, string(data)) - } - } - // do dequeue - for i := 1; i < 10; i++ { - index, bb, err := queue.Dequeue() - if err != nil { - fmt.Println(err) - } - if index != -1 { - fmt.Println(index, string(bb)) - } - - } - - // do gc action to free old data - queue.Gc() - -} - func mainSubscrib() { var queue = new(bigqueue.FileQueue) diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..b47b13b --- /dev/null +++ b/go.mod @@ -0,0 +1,3 @@ +module github.com/jhunters/bigqueue + +go 1.13 diff --git a/mmap_darwin.go b/mmap_darwin.go new file mode 100644 index 0000000..4a4c0fc --- /dev/null +++ b/mmap_darwin.go @@ -0,0 +1,55 @@ +// +build darwin + +package bigqueue + +import ( + "fmt" + "syscall" + "unsafe" +) + +// fdatasync flushes written data to a file descriptor. +func fdatasync(db *DB) error { + return syscall.Fsync(int(db.file.Fd())) +} + +// mmap memory maps a DB's data file. +func mmap(db *DB, sz int) error { + // Map the data file to memory. + b, err := syscall.Mmap(int(db.file.Fd()), 0, sz, syscall.PROT_WRITE|syscall.PROT_READ, syscall.MAP_SHARED|db.MmapFlags) + if err != nil { + return err + } + + // Advise the kernel that the mmap is accessed randomly. + if err := madvise(b, syscall.MADV_RANDOM); err != nil { + return fmt.Errorf("madvise: %s", err) + } + + // Save the original byte slice and convert to a byte array pointer. + db.dataref = b + db.data = (*[maxMapSize]byte)(unsafe.Pointer(&b[0])) + return nil +} + +// munmap unmaps a DB's data file from memory. +func munmap(db *DB) error { + // Ignore the unmap if we have no mapped data. + if db.dataref == nil { + return nil + } + + // Unmap using the original byte slice. + err := syscall.Munmap(db.dataref) + db.data = nil + return err +} + +// NOTE: This function is copied from stdlib because it is not available on darwin. +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 { + err = e1 + } + return +} diff --git a/mmap_unix.go b/mmap_unix.go index f62f449..9024768 100644 --- a/mmap_unix.go +++ b/mmap_unix.go @@ -1,4 +1,4 @@ -// +build !windows,!plan9,!solaris +// +build !windows,!plan9,!solaris,!darwin package bigqueue