diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..6754c7d --- /dev/null +++ b/.gitignore @@ -0,0 +1,2 @@ +cover.out +cover.html diff --git a/.golangci.yaml b/.golangci.yaml new file mode 100644 index 0000000..9aa3e2e --- /dev/null +++ b/.golangci.yaml @@ -0,0 +1,38 @@ +linters: + enable-all: true + disable: + # re-enable when working + - rowserrcheck + - wastedassign + # maybe enable these + - wrapcheck + # leave these disabled + - cyclop + - deadcode + - dupl + - exhaustivestruct + - exhaustruct + - forbidigo + - forcetypeassert + - funlen + - gochecknoglobals + - gocognit + - goconst + - godox + - golint + - gomnd + - ifshort + - interfacer + - lll + - maintidx + - maligned + - nilnil + - nestif + - nlreturn + - nolintlint + - nosnakecase + - scopelint + - structcheck + - thelper + - varcheck + - varnamelen diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..4ec193a --- /dev/null +++ b/go.mod @@ -0,0 +1,21 @@ +module github.com/gopatchy/storebus + +go 1.19 + +require ( + github.com/dchest/uniuri v1.2.0 + github.com/gopatchy/bus v0.0.0-20230420182949-6d46cf96fe01 + github.com/gopatchy/jsrest v0.0.0-20230420161234-12a6d6da8b7f + github.com/gopatchy/metadata v0.0.0-20230420053349-25837551c11d + github.com/gopatchy/store v0.0.0-20230420180842-570716e0aff1 + github.com/stretchr/testify v1.8.2 + go.uber.org/goleak v1.2.1 +) + +require ( + github.com/davecgh/go-spew v1.1.1 // indirect + github.com/mattn/go-sqlite3 v1.14.16 // indirect + github.com/pmezard/go-difflib v1.0.0 // indirect + github.com/vfaronov/httpheader v0.1.0 // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..0f800eb --- /dev/null +++ b/go.sum @@ -0,0 +1,36 @@ +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/dchest/uniuri v1.2.0 h1:koIcOUdrTIivZgSLhHQvKgqdWZq5d7KdMEWF1Ud6+5g= +github.com/dchest/uniuri v1.2.0/go.mod h1:fSzm4SLHzNZvWLvWJew423PhAzkpNQYq+uNLq4kxhkY= +github.com/gopatchy/bus v0.0.0-20230420182949-6d46cf96fe01 h1:GibSdL5yMYle0k6RssUwYJppaRJWaJK5FH4nopUMMsY= +github.com/gopatchy/bus v0.0.0-20230420182949-6d46cf96fe01/go.mod h1:ySA7GifVT/WXU+ZRVF8zj1jW+VUFzKC1iWgV1jzspy0= +github.com/gopatchy/jsrest v0.0.0-20230420161234-12a6d6da8b7f h1:1uGPJm9K0Fro1UEcZpuK6FNPU/U1XX3aS3x0/PdFS40= +github.com/gopatchy/jsrest v0.0.0-20230420161234-12a6d6da8b7f/go.mod h1:Ryi8LRBLFDhQsMQHuh+6VL7HcFWjBXOEiOy9Ip/Q+Ps= +github.com/gopatchy/metadata v0.0.0-20230420053349-25837551c11d h1:chunoM47vkWSanIvLx4uRSkLMG6chDZOy09L2tt/bv8= +github.com/gopatchy/metadata v0.0.0-20230420053349-25837551c11d/go.mod h1:VgD33raUShjDePCDBo55aj+eSXFtUEpMzs+Ie39g2zo= +github.com/gopatchy/store v0.0.0-20230420180842-570716e0aff1 h1:6Van3j5DSGsV5c+OVP7gx2Z3OjAHe9ejlmIAEi0ykJU= +github.com/gopatchy/store v0.0.0-20230420180842-570716e0aff1/go.mod h1:4Z0CHXlhqVJasKMJ8CArEglrZVsqK/k/jj+shCZKJVE= +github.com/kr/pretty v0.3.0 h1:WgNl7dwNpEZ6jJ9k1snq4pZsg7DOEN8hP9Xw0Tsjwk0= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/mattn/go-sqlite3 v1.14.16 h1:yOQRA0RpS5PFz/oikGwBEqvAWhWg5ufRz4ETLjwpU1Y= +github.com/mattn/go-sqlite3 v1.14.16/go.mod h1:2eHXhiwb8IkHr+BDWZGa96P6+rkvnG63S2DGjv9HUNg= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/rogpeppe/go-internal v1.8.1-0.20211023094830-115ce09fd6b4 h1:Ha8xCaq6ln1a+R91Km45Oq6lPXj2Mla6CRJYcuV2h1w= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= +github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo= +github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= +github.com/stretchr/testify v1.8.2 h1:+h33VjcLVPDHtOdpUCuF+7gSuG3yGIftsP1YvFihtJ8= +github.com/stretchr/testify v1.8.2/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= +github.com/vfaronov/httpheader v0.1.0 h1:VdzetvOKRoQVHjSrXcIOwCV6JG5BCAW9rjbVbFPBmb0= +github.com/vfaronov/httpheader v0.1.0/go.mod h1:ZBxgbYu6nbN5V9Ptd1yYUUan0voD0O8nZLXHyxLgoLE= +go.uber.org/goleak v1.2.1 h1:NBol2c7O1ZokfZ0LEU9K6Whx/KnwvepVetCUhtKja4A= +go.uber.org/goleak v1.2.1/go.mod h1:qlT2yGI9QafXHhZZLxlSuNsMw3FFLxBr+tBRlmO1xH4= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/justfile b/justfile new file mode 100644 index 0000000..d607f93 --- /dev/null +++ b/justfile @@ -0,0 +1,18 @@ +go := env_var_or_default('GOCMD', 'go') + +default: tidy test + +tidy: + {{go}} mod tidy + goimports -l -w . + gofumpt -l -w . + {{go}} fmt ./... + +test: + {{go}} vet ./... + golangci-lint run ./... + {{go}} test -race -coverprofile=cover.out -timeout=60s -parallel=10 ./... + {{go}} tool cover -html=cover.out -o=cover.html + +todo: + -git grep -e TODO --and --not -e ignoretodo diff --git a/pkg_test.go b/pkg_test.go new file mode 100644 index 0000000..5ac2997 --- /dev/null +++ b/pkg_test.go @@ -0,0 +1,11 @@ +package storebus_test + +import ( + "testing" + + "go.uber.org/goleak" +) + +func TestMain(m *testing.M) { + goleak.VerifyTestMain(m) +} diff --git a/storebus.go b/storebus.go new file mode 100644 index 0000000..b7c03ec --- /dev/null +++ b/storebus.go @@ -0,0 +1,168 @@ +package storebus + +import ( + "context" + "crypto/sha256" + "encoding/json" + "fmt" + "sync" + + "github.com/gopatchy/bus" + "github.com/gopatchy/jsrest" + "github.com/gopatchy/metadata" + "github.com/gopatchy/store" +) + +type StoreBus struct { + store *store.Store + bus *bus.Bus + + // This lock ensures that no writes interleave with read/subscribe pairs + orderMu sync.RWMutex + + chanMap map[<-chan []any]<-chan any + chanMapMu sync.Mutex +} + +func NewStoreBus(dbname string) (*StoreBus, error) { + st, err := store.NewStore(dbname) + if err != nil { + return nil, err + } + + return &StoreBus{ + store: st, + bus: bus.NewBus(), + chanMap: map[<-chan []any]<-chan any{}, + }, nil +} + +func (sb *StoreBus) Close() { + sb.store.Close() +} + +func (sb *StoreBus) Write(ctx context.Context, t string, obj any) error { + sb.orderMu.Lock() + defer sb.orderMu.Unlock() + + if err := updateHash(obj); err != nil { + return jsrest.Errorf(jsrest.ErrInternalServerError, "hash update failed (%w)", err) + } + + if err := sb.store.Write(ctx, t, obj); err != nil { + return jsrest.Errorf(jsrest.ErrInternalServerError, "write failed (%w)", err) + } + + sb.bus.Announce(t, obj) + + return nil +} + +func (sb *StoreBus) Delete(ctx context.Context, t, id string) error { + sb.orderMu.Lock() + defer sb.orderMu.Unlock() + + if err := sb.store.Delete(ctx, t, id); err != nil { + return jsrest.Errorf(jsrest.ErrInternalServerError, "delete failed (%w)", err) + } + + sb.bus.Delete(t, id) + + return nil +} + +func (sb *StoreBus) Read(ctx context.Context, t, id string, factory func() any) (any, error) { + return sb.store.Read(ctx, t, id, factory) +} + +func (sb *StoreBus) ReadStream(ctx context.Context, t, id string, factory func() any) (<-chan any, error) { + sb.orderMu.RLock() + defer sb.orderMu.RUnlock() + + initial, err := sb.store.Read(ctx, t, id, factory) + if err != nil { + return nil, jsrest.Errorf(jsrest.ErrInternalServerError, "read failed (%w)", err) + } + + c := sb.bus.SubscribeKey(t, id, initial) + + return c, nil +} + +func (sb *StoreBus) CloseReadStream(t, id string, c <-chan any) { + sb.bus.UnsubscribeKey(t, id, c) +} + +func (sb *StoreBus) List(ctx context.Context, t string, factory func() any) ([]any, error) { + return sb.store.List(ctx, t, factory) +} + +func (sb *StoreBus) ListStream(ctx context.Context, t string, factory func() any) (<-chan []any, error) { + sb.orderMu.RLock() + defer sb.orderMu.RUnlock() + + initial, err := sb.store.List(ctx, t, factory) + if err != nil { + return nil, jsrest.Errorf(jsrest.ErrInternalServerError, "list failed (%w)", err) + } + + c := sb.bus.SubscribeType(t, initial) + + ret := make(chan []any, 100) + + sb.registerChan(c, ret) + + go func() { + defer close(ret) + + for range c { + // List() results are always at least (but not exactly) as new as the write that triggered it + l, err := sb.store.List(ctx, t, factory) + if err != nil { + break + } + + select { + case ret <- l: + default: + break + } + } + }() + + return ret, nil +} + +func (sb *StoreBus) CloseListStream(t string, c <-chan []any) { + sb.chanMapMu.Lock() + defer sb.chanMapMu.Unlock() + + sb.bus.UnsubscribeType(t, sb.chanMap[c]) + + delete(sb.chanMap, c) +} + +func (sb *StoreBus) registerChan(in <-chan any, out <-chan []any) { + sb.chanMapMu.Lock() + defer sb.chanMapMu.Unlock() + + sb.chanMap[out] = in +} + +func updateHash(obj any) error { + m := *metadata.GetMetadata(obj) + metadata.ClearMetadata(obj) + + defer metadata.SetMetadata(obj, &m) + + hash := sha256.New() + enc := json.NewEncoder(hash) + + if err := enc.Encode(obj); err != nil { + return jsrest.Errorf(jsrest.ErrInternalServerError, "JSON encode failed (%w)", err) + } + + m.ETag = fmt.Sprintf("etag:%x", hash.Sum(nil)) + + return nil +} diff --git a/storebus_test.go b/storebus_test.go new file mode 100644 index 0000000..94fd5d8 --- /dev/null +++ b/storebus_test.go @@ -0,0 +1,175 @@ +package storebus_test + +import ( + "context" + "fmt" + "os" + "testing" + + "github.com/dchest/uniuri" + "github.com/gopatchy/metadata" + "github.com/gopatchy/storebus" + "github.com/stretchr/testify/require" +) + +func TestStoreBus(t *testing.T) { + t.Parallel() + + dir, err := os.MkdirTemp("", "") + require.NoError(t, err) + + defer os.RemoveAll(dir) + + ctx := context.Background() + + dbname := fmt.Sprintf("file:%s?mode=memory&cache=shared", uniuri.New()) + + sb, err := storebus.NewStoreBus(dbname) + require.NoError(t, err) + + defer sb.Close() + + err = sb.Write(ctx, "storeBusTest", &storeBusTest{ + Metadata: metadata.Metadata{ + ID: "id1", + }, + Opaque: "foo", + }) + require.NoError(t, err) + + c1, err := sb.ReadStream(ctx, "storeBusTest", "id1", newStoreBusTest) + require.NoError(t, err) + + defer sb.CloseReadStream("storeBusTest", "id1", c1) + + c2, err := sb.ListStream(ctx, "storeBusTest", newStoreBusTest) + require.NoError(t, err) + + defer sb.CloseListStream("storeBusTest", c2) + + out1 := (<-c1).(*storeBusTest) + require.Equal(t, "foo", out1.Opaque) + require.Equal(t, "etag:2c8edc6414452b8dee7826bd55e585f850ac47a0dcfc357dc1fcaaa3164cdfa2", out1.ETag) + + l1 := <-c2 + require.Len(t, l1, 1) + require.Equal(t, "foo", l1[0].(*storeBusTest).Opaque) + require.Equal(t, "etag:2c8edc6414452b8dee7826bd55e585f850ac47a0dcfc357dc1fcaaa3164cdfa2", l1[0].(*storeBusTest).ETag) + + err = sb.Write(ctx, "storeBusTest", &storeBusTest{ + Metadata: metadata.Metadata{ + ID: "id1", + }, + Opaque: "bar", + }) + require.NoError(t, err) + + out2 := (<-c1).(*storeBusTest) + require.Equal(t, "bar", out2.Opaque) + require.Equal(t, "etag:906fda69e9893280ca9294bd04eb276794da9a8904fc0b671c69175f08cc03c6", out2.ETag) + + l2 := <-c2 + require.Len(t, l2, 1) + require.Equal(t, "bar", l2[0].(*storeBusTest).Opaque) + require.Equal(t, "etag:906fda69e9893280ca9294bd04eb276794da9a8904fc0b671c69175f08cc03c6", l2[0].(*storeBusTest).ETag) + + l2a, err := sb.List(ctx, "storeBusTest", newStoreBusTest) + require.NoError(t, err) + require.Len(t, l2a, 1) + require.Equal(t, "bar", l2a[0].(*storeBusTest).Opaque) + require.Equal(t, "etag:906fda69e9893280ca9294bd04eb276794da9a8904fc0b671c69175f08cc03c6", l2a[0].(*storeBusTest).ETag) + + out2a, err := sb.Read(ctx, "storeBusTest", "id1", newStoreBusTest) + require.NoError(t, err) + require.Equal(t, "bar", out2a.(*storeBusTest).Opaque) + require.Equal(t, "etag:906fda69e9893280ca9294bd04eb276794da9a8904fc0b671c69175f08cc03c6", out2a.(*storeBusTest).ETag) + + err = sb.Write(ctx, "storeBusTest", &storeBusTest{ + Metadata: metadata.Metadata{ + ID: "id2", + }, + Opaque: "zig", + }) + require.NoError(t, err) + + l3 := <-c2 + require.Len(t, l3, 2) + require.ElementsMatch( + t, + []string{ + l3[0].(*storeBusTest).Opaque, + l3[1].(*storeBusTest).Opaque, + }, + []string{ + "bar", + "zig", + }, + ) + + l3a, err := sb.List(ctx, "storeBusTest", newStoreBusTest) + require.NoError(t, err) + require.Len(t, l3a, 2) + require.ElementsMatch( + t, + []string{ + l3a[0].(*storeBusTest).Opaque, + l3a[1].(*storeBusTest).Opaque, + }, + []string{ + "bar", + "zig", + }, + ) +} + +func TestStoreBusDelete(t *testing.T) { + t.Parallel() + + dir, err := os.MkdirTemp("", "") + require.NoError(t, err) + + defer os.RemoveAll(dir) + + ctx := context.Background() + + dbname := fmt.Sprintf("file:%s?mode=memory&cache=shared", uniuri.New()) + + sb, err := storebus.NewStoreBus(dbname) + require.NoError(t, err) + + defer sb.Close() + + c1, err := sb.ReadStream(ctx, "storeBusTest", "id1", newStoreBusTest) + require.NoError(t, err) + + defer sb.CloseReadStream("storeBusTest", "id1", c1) + + preout := <-c1 + require.Nil(t, preout) + + err = sb.Write(ctx, "storeBusTest", &storeBusTest{ + Metadata: metadata.Metadata{ + ID: "id1", + }, + Opaque: "foo", + }) + require.NoError(t, err) + + out := (<-c1).(*storeBusTest) + require.Equal(t, "foo", out.Opaque) + + err = sb.Delete(ctx, "storeBusTest", "id1") + require.NoError(t, err) + + _, ok := <-c1 + require.False(t, ok) +} + +type storeBusTest struct { + metadata.Metadata + Opaque string +} + +func newStoreBusTest() any { + return &storeBusTest{} +}