refactor: remove dead storm code
7 files changed, 5 insertions(+), 302 deletions(-)
M .golangci.yaml → .golangci.yaml
@@ -19,8 +19,6 @@ - unused settings: errcheck: verbose: true - exclude-functions: - - (github.com/asdine/storm/v3.Tx).Rollback forbidigo: forbid: - pattern: ^filepath\.(Join|Walkdir)$
M go.mod → go.mod
@@ -7,7 +7,6 @@ alin.ovh/gomponents v1.6.0 alin.ovh/x v1.0.0 github.com/Southclaws/fault v0.8.2 github.com/andybalholm/brotli v1.1.1 - github.com/asdine/storm/v3 v3.2.1 github.com/bcicen/jstream v1.0.1 github.com/blevesearch/bleve/v2 v2.5.2 github.com/blevesearch/bleve_index_api v1.2.8@@ -24,7 +23,6 @@ github.com/pelletier/go-toml/v2 v2.2.4 github.com/stefanfritsch/goldmark-fences v1.0.0 github.com/stoewer/go-strcase v1.3.0 github.com/yuin/goldmark v1.7.12 - go.etcd.io/bbolt v1.4.2 go.uber.org/zap v1.27.0 golang.org/x/net v0.41.0 golang.org/x/term v0.34.0@@ -64,6 +62,7 @@ github.com/pkg/errors v0.9.1 // indirect github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect github.com/sykesm/zap-logfmt v0.0.4 // indirect github.com/thessem/zap-prettyconsole v0.5.2 // indirect + go.etcd.io/bbolt v1.4.2 // indirect go.uber.org/multierr v1.11.0 // indirect golang.org/x/exp v0.0.0-20250606033433-dcc06ee1d476 // indirect golang.org/x/sys v0.35.0 // indirect
M go.sum → go.sum
@@ -5,18 +5,12 @@ alin.ovh/x v1.0.0/go.mod h1:ivPzLVkiYeXGmUITqx4KLoNikO1URKJXYPPo4yCxuhY= github.com/BurntSushi/toml v0.3.1/go.mod h1:xHWCNGjB5oqiDr8zfno3MHue2Ht5sIBksp03qcyfWMU= github.com/Code-Hex/dd v1.1.0 h1:VEtTThnS9l7WhpKUIpdcWaf0B8Vp0LeeSEsxA1DZseI= github.com/Code-Hex/dd v1.1.0/go.mod h1:VaMyo/YjTJ3d4qm/bgtrUkT2w+aYwJ07Y7eCWyrJr1w= -github.com/DataDog/zstd v1.4.1 h1:3oxKN3wbHibqx897utPC2LTQU4J+IHWWJO+glkAkpFM= -github.com/DataDog/zstd v1.4.1/go.mod h1:1jcaCB/ufaK+sKp1NBhlGmpz41jOoPQ35bpF36t7BBo= github.com/RoaringBitmap/roaring/v2 v2.5.0 h1:TJ45qCM7D7fIEBwKd9zhoR0/S1egfnSSIzLU1e1eYLY= github.com/RoaringBitmap/roaring/v2 v2.5.0/go.mod h1:FiJcsfkGje/nZBZgCu0ZxCPOKD/hVXDS2dXi7/eUFE0= -github.com/Sereal/Sereal v0.0.0-20190618215532-0b8ac451a863 h1:BRrxwOZBolJN4gIwvZMJY1tzqBvQgpaZiQRuIDD40jM= -github.com/Sereal/Sereal v0.0.0-20190618215532-0b8ac451a863/go.mod h1:D0JMgToj/WdxCgd30Kc1UcA9E+WdZoJqeVOuYW7iTBM= github.com/Southclaws/fault v0.8.2 h1:hbQANoRWYVWnQjpwJlNlfaolM+oIihgoFowaY3EBLCs= github.com/Southclaws/fault v0.8.2/go.mod h1:VUVkAWutC59SL16s6FTqf3I6I2z77RmnaW5XRz4bLOE= github.com/andybalholm/brotli v1.1.1 h1:PR2pgnyFznKEugtsUo0xLdDop5SKXd5Qf5ysW+7XdTA= github.com/andybalholm/brotli v1.1.1/go.mod h1:05ib4cKhjx3OQYUY22hTVd34Bc8upXjOLL2rKwwZBoA= -github.com/asdine/storm/v3 v3.2.1 h1:I5AqhkPK6nBZ/qJXySdI7ot5BlXSZ7qvDY1zAn5ZJac= -github.com/asdine/storm/v3 v3.2.1/go.mod h1:LEpXwGt4pIqrE/XcTvCnZHT5MgZCV6Ub9q7yQzOFWr0= github.com/bcicen/jstream v1.0.1 h1:BXY7Cu4rdmc0rhyTVyT3UkxAiX3bnLpKLas9btbH5ck= github.com/bcicen/jstream v1.0.1/go.mod h1:9ielPxqFry7Y4Tg3j4BfjPocfJ3TbsRtXOAYXYmRuAQ= github.com/benbjohnson/clock v1.1.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA=@@ -77,11 +71,8 @@ github.com/getsentry/sentry-go v0.33.0 h1:YWyDii0KGVov3xOaamOnF0mjOrqSjBqwv48UEzn7QFg= github.com/getsentry/sentry-go v0.33.0/go.mod h1:C55omcY9ChRQIUcVcGcs+Zdy4ZpQGvNJ7JYHIoSWOtE= github.com/go-errors/errors v1.4.2 h1:J6MZopCL4uSllY1OfXM374weqZFFItUbrImctkmUxIA= github.com/go-errors/errors v1.4.2/go.mod h1:sIVyrIiJhuEF+Pj9Ebtd6P/rEYROXFi3BopGUQ5a5Og= -github.com/golang/protobuf v1.3.1/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= -github.com/golang/protobuf v1.3.2/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= -github.com/golang/snappy v0.0.1/go.mod h1:/XxbfmMg8lxefKM7IXC3fBNl/7bRcc72aCRzEWrmP2Q= github.com/golang/snappy v1.0.0 h1:Oy607GVXHs7RtbggtPBnr2RmDArIsAefDwvrdWvRhGs= github.com/golang/snappy v1.0.0/go.mod h1:/XxbfmMg8lxefKM7IXC3fBNl/7bRcc72aCRzEWrmP2Q= github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8=@@ -136,7 +127,6 @@ github.com/stoewer/go-strcase v1.3.0/go.mod h1:fAH5hQ5pehh+j3nZfvwdk2RgEgQjAoM8wodgtPmh1xo= 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.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXfy6kDkUVs= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=@@ -150,15 +140,12 @@ github.com/sykesm/zap-logfmt v0.0.4/go.mod h1:AuBd9xQjAe3URrWT1BBDk2v2onAZHkZkWRMiYZXiZWA= github.com/tailscale/depaware v0.0.0-20210622194025-720c4b409502/go.mod h1:p9lPsd+cx33L3H9nNoecRRxPssFKUwwI50I3pZ0yT+8= github.com/thessem/zap-prettyconsole v0.5.2 h1:knusxXGhmkD5Ho+WiI4IzD16Dz9PEcOIKdK+uX4oTPA= github.com/thessem/zap-prettyconsole v0.5.2/go.mod h1:3qfsE7y+bLOq7EQ+fMZHD3HYEp24ULFf5nhLSx6rjrE= -github.com/vmihailenco/msgpack v4.0.4+incompatible h1:dSLoQfGFAo3F6OoNhwUmLwVgaUXK79GlxNBwueZn0xI= -github.com/vmihailenco/msgpack v4.0.4+incompatible/go.mod h1:fy3FlTQTDXWkZ7Bh6AcGMlsjHatGryHQYUTf1ShIgkk= github.com/xyproto/randomstring v1.0.5 h1:YtlWPoRdgMu3NZtP45drfy1GKoojuR7hmRcnhZqKjWU= github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E= github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.3.5/go.mod h1:mwnBkeHKe2W/ZEtQ+71ViKU8L12m81fl3OWwC1Zlc8k= github.com/yuin/goldmark v1.7.12 h1:YwGP/rrea2/CnCtUHgjuolG/PnMxdQtPMO5PvaE2/nY= github.com/yuin/goldmark v1.7.12/go.mod h1:ip/1k0VRfGynBgxOz0yCqHrbZXhcjxyuS66Brc7iBKg= -go.etcd.io/bbolt v1.3.4/go.mod h1:G5EMThwa9y8QZGBClrRx5EY+Yw9kAhnjy3bSjsnlVTQ= go.etcd.io/bbolt v1.4.2 h1:IrUHp260R8c+zYx/Tm8QZr04CX+qWS5PGfPdevhdm1I= go.etcd.io/bbolt v1.4.2/go.mod h1:Is8rSHO/b4f3XigBC0lL0+4FwAQv3HXEEIgFMuKHceM= go.uber.org/atomic v1.5.0/go.mod h1:sABNBOSYdrvTF6hTgEIbc7YasKWGhgEQZyfxyTvoXHQ=@@ -192,9 +179,7 @@ golang.org/x/mod v0.25.0 h1:n7a+ZbQKQA/Ysbyb0/6IbB1H/X41mKgbhfv7AfG/44w= golang.org/x/mod v0.25.0/go.mod h1:IXM97Txy2VM4PJ3gI61r1YEk/gAj6zAHN3AdZt6S9Ww= golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= -golang.org/x/net v0.0.0-20190603091049-60506f45cf65/go.mod h1:HSz+uSET+XFnRR8LxR5pz3Of3rY3CfYBVs4xY44aLks= golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= -golang.org/x/net v0.0.0-20191105084925-a882066a44e0/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20201021035429-f5854403a974/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU= golang.org/x/net v0.0.0-20210405180319-a5a99cb37ef4/go.mod h1:p54w0d4576C0XHj96bSt6lcn1PtDYWL6XObtHCRCNQM= golang.org/x/net v0.41.0 h1:vBTly1HeNPEn3wtREYfy4GZ/NECgw2Cnl+nK6Nz3uvw=@@ -206,7 +191,6 @@ golang.org/x/sync v0.15.0 h1:KWH3jNZsfyT6xfAfKiz6MRNmd46ByHDYaZ7KSkCtdW8= golang.org/x/sync v0.15.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20200202164722-d101bd2416d5/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20200930185726-fdedc70b468f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210330210617-4fbd30eecc44/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=@@ -221,7 +205,6 @@ golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/term v0.34.0 h1:O/2T7POpk0ZZ7MAzMeWFSg6S5IpWd/RXDlM9hgM3DR4= golang.org/x/term v0.34.0/go.mod h1:5jC53AEywhIVebHgPVeg0mj8OD3VO9OzclacVrqpaAw= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= -golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk= golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.26.0 h1:P42AVeLghgTYr4+xUnTRKDMqpar+PtX7KWuNQL21L8M= golang.org/x/text v0.26.0/go.mod h1:QK15LZJUUQVJxhz7wXgxSy/CJaTFjd0G+YLonydOVQA=@@ -238,13 +221,10 @@ golang.org/x/tools v0.34.0/go.mod h1:pAP9OwEaY1CAW3HOmg3hLZC5Z0CCmzjAF2UQMSqNARg= golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -google.golang.org/appengine v1.6.5 h1:tycE03LOZYQNhDpS27tcQdAzLCVMaj7QT2SXxebnpCM= -google.golang.org/appengine v1.6.5/go.mod h1:8WjMMxjGQR8xUklV/ARdw2HLXBOI7O7uCIDZVag1xfc= google.golang.org/protobuf v1.36.6 h1:z1NpPI8ku2WgiWnf+t9wTPsn6eP1L7ksHUlkfLvd9xY= google.golang.org/protobuf v1.36.6/go.mod h1:jduwjTPXsFjZGTmRluh+L6NjiWu7pchiJ2/5YcXBHnY= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/errgo.v2 v2.1.0/go.mod h1:hNsd1EY+bozCKY1Ytp96fpM3vjJbqLJn88ws8XvfDNI= gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= gopkg.in/yaml.v2 v2.2.8/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
M gomod2nix.toml → gomod2nix.toml
@@ -19,9 +19,6 @@ hash = "sha256-ifMHIRqeMNV2+jhi5jxA42iK1DATNPYI4pN69inRGT4=" [mod."github.com/andybalholm/brotli"] version = "v1.1.1" hash = "sha256-kCt+irK1gvz2lGQUeEolYa5+FbLsfWlJMCd5hm+RPgQ=" - [mod."github.com/asdine/storm/v3"] - version = "v3.2.1" - hash = "sha256-BLpBFWFjLd5Xumx72cgSA2Zget+4UYMTW+z4bkkdIR0=" [mod."github.com/bcicen/jstream"] version = "v1.0.1" hash = "sha256-mm+/BuIEYYj6XOHCCJLxVMKd1XcBXCiRCWA+aTvr1sE="
M internal/nix/option.go → internal/nix/option.go
@@ -23,7 +23,7 @@ URL string } type Option struct { - Name string `storm:"id"` + Name string Source string Declarations []Link Default *Docs `json:",omitempty"`@@ -33,7 +33,7 @@ Loc []string Parents string RelatedPackages Markdown `json:",omitempty"` Type string - ImportedAt time.Time `storm:"index"` + ImportedAt time.Time } func (Option) BleveType() string {
M internal/nix/package.go → internal/nix/package.go
@@ -4,7 +4,7 @@ import "time" type Package struct { Name string - Attribute string `storm:"id"` + Attribute string Source string Broken bool Definition string@@ -18,7 +18,7 @@ Maintainers []Maintainer PackageSet string Platforms []string Version string - ImportedAt time.Time `storm:"index"` + ImportedAt time.Time } type License struct {
D internal/storage/store.go
@@ -1,271 +0,0 @@ -package storage - -import ( - "context" - "errors" - "time" - - "alin.ovh/x/log" - "github.com/Southclaws/fault" - "github.com/Southclaws/fault/fmsg" - "github.com/Southclaws/fault/ftag" - "github.com/asdine/storm/v3" - "github.com/asdine/storm/v3/codec/gob" - "go.etcd.io/bbolt" - - "alin.ovh/searchix/internal/config" - "alin.ovh/searchix/internal/file" - "alin.ovh/searchix/internal/nix" -) - -var BatchSize = 50000 - -type Options struct { - Replace bool - LowMemory bool - Root *file.Root - Logger *log.Logger -} - -type Store struct { - *storm.DB - new bool - log *log.Logger -} - -const filename = "searchix.bolt" - -func New(opts *Options) (*Store, error) { - exists, err := opts.Root.Exists(filename) - if err != nil { - return nil, fault.Wrap(err, fmsg.With("failed to check if file exists")) - } - - if opts.Replace && exists { - err = opts.Root.Remove(filename) - if err != nil { - return nil, fault.Wrap(err, fmsg.With("failed to remove existing file")) - } - - exists = false - } - - //nolint:forbidigo // external package - path := opts.Root.JoinPath(filename) - bb, err := storm.Open(path, - storm.Codec(gob.Codec), - storm.BoltOptions(0o600, &bbolt.Options{ - FreelistType: bbolt.FreelistMapType, - NoFreelistSync: true, - NoGrowSync: true, - Timeout: 1 * time.Second, - }), - ) - if err != nil { - return nil, fault.Wrap(err, fmsg.With("failed to open database")) - } - - if !opts.LowMemory { - bb.Bolt.AllocSize = 256 * 1024 * 1024 - } - - return &Store{ - DB: bb, - new: !exists, - log: opts.Logger, - }, nil -} - -func (s *Store) IsNew() bool { - return s.new -} - -func (s *Store) Close() error { - err := s.DB.Close() - if err != nil { - return fault.Wrap(err, fmsg.With("failed to close database")) - } - - return nil -} - -func (s *Store) MakeSourceImporter( - source *config.Source, -) func(context.Context, <-chan nix.Importable) <-chan error { - return func(ctx context.Context, objects <-chan nix.Importable) <-chan error { - var imp nix.Importable - - errs := make(chan error) - node := s.DB.From(source.Key).WithBatch(true) - - i := 0 - - var save func(storm.Node, nix.Importable) error - switch source.Importer { - case config.Packages: - imp = &nix.Package{} - save = saveGen[nix.Package] - case config.Options: - imp = &nix.Option{} - save = saveGen[nix.Option] - default: - errs <- fault.New("invalid importer") - - return errs - } - - go func() { - defer close(errs) - tx, err := node.Begin(true) - if err != nil { - errs <- fault.Wrap(err, fmsg.With("failed to begin transaction")) - - return - } - defer func() { - if err := tx.Rollback(); err != nil { - if !errors.Is(err, storm.ErrNotInTransaction) { - errs <- fault.Wrap(err, fmsg.With("failed to rollback transaction")) - } - } - }() - - outer: - for obj := range objects { - i++ - select { - case <-ctx.Done(): - s.log.Warn("import aborted") - - break outer - default: - } - - err := save(tx, obj) - if err != nil { - errs <- fault.Wrap(err, fmsg.With("failed to save object")) - } - - if i%BatchSize == 0 { - s.log.Debug("imported", "count", i) - err := tx.Commit() - if err != nil { - errs <- fault.Wrap(err, fmsg.With("failed to commit transaction")) - } - tx, err = node.Begin(true) - if err != nil { - errs <- fault.Wrap(err, fmsg.With("failed to begin transaction")) - } - } - } - - if err := tx.Commit(); err != nil { - errs <- fault.Wrap(err, fmsg.With("failed to commit transaction")) - } - - if err := node.ReIndex(imp); err != nil { - errs <- fault.Wrap(err, fmsg.With("failed to reindex storm db")) - } - }() - - return errs - } -} - -func saveGen[T nix.Importable](node storm.Node, obj nix.Importable) error { - doc, ok := obj.(T) - if !ok { - return fault.Newf("invalid type: %T", obj) - } - - if err := node.Save(&doc); err != nil { - return fault.Wrap(err, fmsg.With("failed to save document")) - } - - return nil -} - -func (s *Store) GetDocument( - source *config.Source, - id string, -) (nix.Importable, error) { - var doc nix.Importable - var err error - - node := s.From(source.Key) - - switch source.Importer { - case config.Packages: - doc = &nix.Package{} - err = node.One("Attribute", id, doc) - case config.Options: - doc = &nix.Option{} - err = node.One("Name", id, doc) - default: - return nil, fault.New("invalid importer type") - } - - if err != nil { - if errors.Is(err, storm.ErrNotFound) { - return nil, - fault.Wrap( - fault.Newf("document not found source: %s id: %s", source.Key, id), - ftag.With(ftag.NotFound), - ) - } - - return nil, fault.Wrap( - err, - fmsg.Withf("failed to get document source: %s id: %s", source.Key, id), - ) - } - - return doc, nil -} - -func MakeSourceExporter[T nix.Importable]( - store *Store, - source *config.Source, - batchSize int, -) func(context.Context) (<-chan nix.Importable, <-chan error) { - return func(_ context.Context) (<-chan nix.Importable, <-chan error) { - results := make(chan nix.Importable, 1) - errs := make(chan error) - - go func() { - defer close(results) - defer close(errs) - - var obj T - objs := make([]T, 0, batchSize) - node := store.From(source.Key) - count, err := node.Count(&obj) - if err != nil { - errs <- fault.Wrap( - err, - fmsg.Withf("failed to count documents source: %s", source.Key), - ) - - return - } - - limit := min(batchSize, count) - for offset := 0; offset < count; offset += batchSize { - err := node.All(&objs, storm.Skip(offset), storm.Limit(limit)) - if err != nil { - errs <- fault.Wrap( - err, - fmsg.Withf("failed to export documents source: %s offset: %d limit: %d", source.Key, offset, limit), - ) - - return - } - for _, obj := range objs { - results <- obj - } - } - }() - - return results, errs - } -}