diff --git a/go.mod b/go.mod index 482865f8..182e5ee0 100644 --- a/go.mod +++ b/go.mod @@ -3,7 +3,6 @@ module oj go 1.25.0 require ( - github.com/acaloiaro/neoq v0.19.0 github.com/alexandrevicenzi/go-sse v1.6.0 github.com/georgysavva/scany/v2 v2.1.4 github.com/go-chi/chi/v5 v5.2.4 @@ -53,9 +52,7 @@ require ( github.com/google/uuid v1.6.0 // indirect github.com/gorilla/css v1.0.0 // indirect github.com/gorilla/mux v1.8.0 // indirect - github.com/guregu/null v4.0.0+incompatible // indirect github.com/huandu/xstrings v1.5.0 // indirect - github.com/iancoleman/strcase v0.3.0 // indirect github.com/inconshreveable/mousetrap v1.1.0 // indirect github.com/jackc/pgerrcode v0.0.0-20250907135507-afb5586c32a6 // indirect github.com/jackc/pgpassfile v1.0.0 // indirect @@ -63,7 +60,6 @@ require ( github.com/jackc/puddle/v2 v2.2.2 // indirect github.com/jinzhu/inflection v1.0.0 // indirect github.com/json-iterator/go v1.1.12 // indirect - github.com/jsuar/go-cron-descriptor v0.1.0 // indirect github.com/klauspost/compress v1.16.7 // indirect github.com/klauspost/cpuid/v2 v2.2.5 // indirect github.com/leodido/go-urn v1.4.0 // indirect @@ -88,7 +84,6 @@ require ( github.com/riverqueue/river/rivershared v0.39.0 // indirect github.com/riverqueue/river/rivertype v0.39.0 // indirect github.com/riza-io/grpc-go v0.2.0 // indirect - github.com/robfig/cron v1.2.0 // indirect github.com/rs/xid v1.5.0 // indirect github.com/shopspring/decimal v1.4.0 // indirect github.com/sirupsen/logrus v1.9.3 // indirect diff --git a/go.sum b/go.sum index ff86d892..af0a6c0d 100644 --- a/go.sum +++ b/go.sum @@ -15,8 +15,6 @@ github.com/Microsoft/go-winio v0.6.2 h1:F2VQgta7ecxGYO8k3ZZz3RS8fVIXVxONVUPlNERo github.com/Microsoft/go-winio v0.6.2/go.mod h1:yd8OoFMLzJbo9gZq8j5qaps8bJ9aShtEA8Ipt1oGCvU= github.com/ProtonMail/go-crypto v1.3.0 h1:ILq8+Sf5If5DCpHQp4PbZdS1J7HDFRXz/+xKBiRGFrw= github.com/ProtonMail/go-crypto v1.3.0/go.mod h1:9whxjD8Rbs29b4XWbB8irEcE8KHMqaR2e7GWU1R+/PE= -github.com/acaloiaro/neoq v0.19.0 h1:/icF/sDogx7Hi3s6ACtlnGPfHmlxps4SE1RgVb0mz4Q= -github.com/acaloiaro/neoq v0.19.0/go.mod h1:uLI2JfC2V9p+H5lx5s0SUsjxIYLm8JDhi6VkyURMvLM= github.com/ajstarks/svgo v0.0.0-20200320125537-f189e35d30ca/go.mod h1:K08gAheRH3/J6wwsYMMT4xOr94bZjxIelGM0+d/wbFw= github.com/alexandrevicenzi/go-sse v1.6.0 h1:3KvOzpuY7UrbqZgAtOEmub9/V5ykr7Myudw+PA+H1Ik= github.com/alexandrevicenzi/go-sse v1.6.0/go.mod h1:jdrNAhMgVqP7OfcUuM8eJx0sOY17wc+girs5utpFZUU= @@ -27,7 +25,6 @@ github.com/aymerick/douceur v0.2.0/go.mod h1:wlT5vV2O3h55X9m7iVYN0TBM0NH/MmbLnd3 github.com/benbjohnson/clock v1.1.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= -github.com/client9/misspell v0.3.4/go.mod h1:qj6jICC3Q7zFZvVWo7KLAzC3yx5G7kyvSDkc90ppPyw= github.com/cloudflare/circl v1.6.3 h1:9GPOhQGF9MCYUeXyMYlqTR6a5gTrgR/fBLXvUgtVcg8= github.com/cloudflare/circl v1.6.3/go.mod h1:2eXP6Qfat4O/Yhh8BznvKnJ+uzEoTQ6jVKJRn81BiS4= github.com/cockroachdb/cockroach-go/v2 v2.2.0 h1:/5znzg5n373N/3ESjHF5SMLxiW4RKB05Ql//KWfeTFs= @@ -84,21 +81,16 @@ github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= github.com/google/pprof v0.0.0-20250317173921-a4b03ec1a45e h1:ijClszYn+mADRFY17kjQEVQ1XRhq2/JR1M3sGqeJoxs= github.com/google/pprof v0.0.0-20250317173921-a4b03ec1a45e/go.mod h1:boTsfXsheKC2y+lKOCMpSfarhxDeIzfZG1jqGcPl3cA= -github.com/google/renameio v0.1.0/go.mod h1:KWCgfxg9yswjAJkECMjeO8J8rahYeXnNhOm40UhjYkI= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/gorilla/css v1.0.0 h1:BQqNyPTi50JCFMTw/b67hByjMVXZRwGha6wxVGkeihY= github.com/gorilla/css v1.0.0/go.mod h1:Dn721qIggHpt4+EFCcTLTU/vk5ySda2ReITrtgBl60c= github.com/gorilla/mux v1.8.0 h1:i40aqfkR1h2SlN9hojwV5ZA91wcXFOvkdNIeFDP5koI= github.com/gorilla/mux v1.8.0/go.mod h1:DVbg23sWSpFRCP0SfiEN6jmj59UnW/n46BH5rLB71So= -github.com/guregu/null v4.0.0+incompatible h1:4zw0ckM7ECd6FNNddc3Fu4aty9nTlpkkzH7dPn4/4Gw= -github.com/guregu/null v4.0.0+incompatible/go.mod h1:ePGpQaN9cw0tj45IR5E5ehMvsFlLlQZAkkOXZurJ3NM= github.com/hako/durafmt v0.0.0-20210608085754-5c1018a4e16b h1:wDUNC2eKiL35DbLvsDhiblTUXHxcOPwQSCzi7xpQUN4= github.com/hako/durafmt v0.0.0-20210608085754-5c1018a4e16b/go.mod h1:VzxiSdG6j1pi7rwGm/xYI5RbtpBgM8sARDXlvEvxlu0= github.com/huandu/xstrings v1.5.0 h1:2ag3IFq9ZDANvthTwTiqSSZLjDc+BedvHPAp5tJy2TI= github.com/huandu/xstrings v1.5.0/go.mod h1:y5/lhBue+AyNmUVz9RLU9xbLR0o4KIIExikq4ovT0aE= -github.com/iancoleman/strcase v0.3.0 h1:nTXanmYxhfFAMjZL34Ov6gkzEsSJZ5DbhxWjvSASxEI= -github.com/iancoleman/strcase v0.3.0/go.mod h1:iwCmte+B7n89clKwxIoIXy/HfoL7AsD47ZCWhYzw7ho= github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8= github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw= github.com/jackc/pgerrcode v0.0.0-20250907135507-afb5586c32a6 h1:D/V0gu4zQ3cL2WKeVNVM4r2gLxGGf6McLwgXzRTo2RQ= @@ -120,9 +112,6 @@ github.com/jmoiron/sqlx v1.3.5/go.mod h1:nRVWtLre0KfCLJvgxzCsLVMogSvQ1zNJtpYr2Cc github.com/json-iterator/go v1.1.10/go.mod h1:KdQUCv79m/52Kvf8AW2vK1V8akMuk1QjK/uOdHXbAo4= github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= -github.com/jsuar/go-cron-descriptor v0.1.0 h1:Q97ujk+/xhcz1lA9nmUMq750FZV7RlV23TNyMB74xkQ= -github.com/jsuar/go-cron-descriptor v0.1.0/go.mod h1:PFR+Y6Lr86uYZpwsWRoFMRA3CX4a6q7zfPltp3SLkSU= -github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= github.com/klauspost/compress v1.16.7 h1:2mk3MPGNzKyxErAw8YaohYh69+pa4sIQSC0fPGCFR9I= github.com/klauspost/compress v1.16.7/go.mod h1:ntbaceVETuRiXiv4DpjP66DpAtAGkEQskQzEyD//IeE= github.com/klauspost/cpuid/v2 v2.0.1/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg= @@ -206,11 +195,8 @@ github.com/riverqueue/river/rivertype v0.39.0 h1:0jHUTRDR1kdzbgXc6lN1B93WxolZyqP github.com/riverqueue/river/rivertype v0.39.0/go.mod h1:D1Ad+EaZiaXbQbJcJcfeicXJMBKno0n6UcfKI5Q7DIQ= github.com/riza-io/grpc-go v0.2.0 h1:2HxQKFVE7VuYstcJ8zqpN84VnAoJ4dCL6YFhJewNcHQ= github.com/riza-io/grpc-go v0.2.0/go.mod h1:2bDvR9KkKC3KhtlSHfR3dAXjUMT86kg4UfWFyVGWqi8= -github.com/robfig/cron v1.2.0 h1:ZjScXvvxeQ63Dbyxy76Fj3AT3Ut0aKsyd2/tl3DTMuQ= -github.com/robfig/cron v1.2.0/go.mod h1:JGuDeoQd7Z6yL4zQhZ3OPEVHB7fL6Ka6skscFHfmt2k= github.com/robfig/cron/v3 v3.0.1 h1:WdRxkvbJztn8LMz/QEvLN5sBU+xKpSqwwUO1Pjr4qDs= github.com/robfig/cron/v3 v3.0.1/go.mod h1:eQICP3HwyT7UooqI/z+Ov+PtYAWygg1TEWWzGIFLtro= -github.com/rogpeppe/go-internal v1.3.0/go.mod h1:M8bDsm7K2OlrFYOpmOWEs/qY81heoFRclV5y23lUDJ4= github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= github.com/rs/xid v1.5.0 h1:mKX4bl4iPYJtEIxp6CYiUuLQ/8DYMoz0PUdtGgMFRVc= @@ -284,28 +270,22 @@ go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0 go.uber.org/goleak v1.1.10/go.mod h1:8a7PlsEVH3e/a/GLqe5IIrQx6GzcnRmZEufDUTk4A7A= go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= -go.uber.org/multierr v1.5.0/go.mod h1:FeouvMocqHpRaaGuG9EjoKcStLC43Zu/fmqdUMPcKYU= go.uber.org/multierr v1.6.0/go.mod h1:cdWPpRnG4AhwMwsgIHip0KRBQjJy5kYEpYjJxpXp9iU= go.uber.org/multierr v1.7.0/go.mod h1:7EAYxJLBy9rStEaz58O2t4Uvip6FSURkq8/ppBp95ak= go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y= -go.uber.org/tools v0.0.0-20190618225709-2cfd321de3ee/go.mod h1:vJERXedbb3MVM5f9Ejo0C68/HhF8uaILCdgjnY+goOA= -go.uber.org/zap v1.15.0/go.mod h1:Mb2vm2krFEG5DV0W9qcHBYFtp/Wku1cvYaqPsS/WYfc= go.uber.org/zap v1.19.0/go.mod h1:xg/QME4nWcxGxrpdeYfq7UvYrLh66cuVKdrbD1XF/NI= go.uber.org/zap v1.27.0 h1:aJMhYGrd5QSmlpLMr2MftRKl7t8J8PTZPA732ud/XR8= go.uber.org/zap v1.27.0/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= -golang.org/x/crypto v0.0.0-20190510104115-cbcb75029529/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.48.0 h1:/VRzVqiRSggnhY7gNRxPauEQ5Drw9haKdM0jqfcCFts= golang.org/x/crypto v0.48.0/go.mod h1:r0kV5h3qnFPlQnBSrULhlsRfryS2pmewsg+XfMgkVos= golang.org/x/exp v0.0.0-20250305212735-054e65f0b394 h1:nDVHiLt8aIbd/VzvPWN6kSOPE7+F/fNFDSXLVYkE/Iw= golang.org/x/exp v0.0.0-20250305212735-054e65f0b394/go.mod h1:sIifuuw/Yco/y6yb6+bDNfyeQ/MdPUy/hKEMYQV17cM= golang.org/x/lint v0.0.0-20190930215403-16217165b5de/go.mod h1:6SW0HCj/g11FgYtHlgUYUwCkIfeOF89ocIRzGO/8vkc= -golang.org/x/mod v0.0.0-20190513183733-4bf6d317e70e/go.mod h1:mXi4GBBbnImb6dmsKGUJ2LatrhH/nqhxcFungHvyanc= golang.org/x/mod v0.36.0 h1:JJjpVx6myfUsUdAzZuOSTTmRE0PfZeNWzzvKrP7amb4= golang.org/x/mod v0.36.0/go.mod h1:moc6ELqsWcOw5Ef3xVprK5ul/MvtVvkIXLziUOICjUQ= 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-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.49.0 h1:eeHFmOGUTtaaPSGNmjBKpbng9MulQsJURQUAfUwY++o= golang.org/x/net v0.49.0/go.mod h1:/ysNB2EvaqvesRkuLAyjI1ycPZlQHM3q01F02UY/MV8= @@ -313,7 +293,6 @@ golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJ golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= 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-20220715151400-c0bba94af5f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= @@ -325,9 +304,7 @@ golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc= golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38= golang.org/x/tools v0.0.0-20190311212946-11955173bddd/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs= -golang.org/x/tools v0.0.0-20190621195816-6e04913cbbac/go.mod h1:/rFqwRUd4F7ZHNgwSSTFct+R/Kf4OFW1sUzUTQQTgfc= golang.org/x/tools v0.0.0-20191029041327-9cc4af7d6b2c/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= -golang.org/x/tools v0.0.0-20191029190741-b9c20aec41a5/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= golang.org/x/tools v0.0.0-20191108193012-7d206e10da11/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= golang.org/x/tools v0.44.0 h1:UP4ajHPIcuMjT1GqzDWRlalUEoY+uzoZKnhOjbIPD2c= golang.org/x/tools v0.44.0/go.mod h1:KA0AfVErSdxRZIsOVipbv3rQhVXTnlU6UhKxHd1seDI= @@ -349,7 +326,6 @@ gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8 gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= -gopkg.in/errgo.v2 v2.1.0/go.mod h1:hNsd1EY+bozCKY1Ytp96fpM3vjJbqLJn88ws8XvfDNI= gopkg.in/ini.v1 v1.67.0 h1:Dgnx+6+nfE+IfzjUEISNeydPJh9AXNNsWbGP9KzCsOA= gopkg.in/ini.v1 v1.67.0/go.mod h1:pNLf8WUiyNEtQjuu5G5vTm06TEv9tsIgeAvK8hOrP4k= gopkg.in/natefinch/lumberjack.v2 v2.0.0/go.mod h1:l0ndWWf7gzL7RNwBG7wST/UCcT4T24xpD6X8LsfU/+k= @@ -361,7 +337,6 @@ gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C gopkg.in/yaml.v3 v3.0.0-20210107192922-496545a6307b/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= -honnef.co/go/tools v0.0.1-2019.2.3/go.mod h1:a3bituU0lyd329TUQxRnasdCoJDkEUEAqEt0JzvZhAg= maragu.dev/gomponents v1.1.0 h1:iCybZZChHr1eSlvkWp/JP3CrZGzctLudQ/JI3sBcO4U= maragu.dev/gomponents v1.1.0/go.mod h1:oEDahza2gZoXDoDHhw8jBNgH+3UR5ni7Ur648HORydM= modernc.org/cc/v4 v4.25.2 h1:T2oH7sZdGvTaie0BRNFbIYsabzCxUQg8nLqCdQ2i0ic= diff --git a/worker/notifydelivery/notifydelivery.go b/worker/notifydelivery/notifydelivery.go index d13b782c..5ce72586 100644 --- a/worker/notifydelivery/notifydelivery.go +++ b/worker/notifydelivery/notifydelivery.go @@ -10,30 +10,29 @@ import ( "oj/services/email" "time" - "github.com/acaloiaro/neoq/jobs" "github.com/georgysavva/scany/v2/pgxscan" "github.com/jackc/pgx/v5/pgxpool" + "github.com/riverqueue/river" ) -type service struct { - Queries *api.Queries - Conn *pgxpool.Pool +type NotifyDeliveryArgs struct { + ID int64 `json:"id"` } -func NewService(q *api.Queries, conn *pgxpool.Pool) *service { - return &service{Queries: q, Conn: conn} +func (NotifyDeliveryArgs) Kind() string { return "notify_delivery" } + +type Worker struct { + river.WorkerDefaults[NotifyDeliveryArgs] + Queries *api.Queries + Conn *pgxpool.Pool } -func (s *service) Handle(ctx context.Context) error { - j, err := jobs.FromContext(ctx) - if err != nil { - return err - } - log.Printf("handleNotifyDelivery job id: %d, payload: %v", j.ID, j.Payload) +func (w *Worker) Work(ctx context.Context, job *river.Job[NotifyDeliveryArgs]) error { + log.Printf("handleNotifyDelivery job id: %d, delivery id: %d", job.ID, job.Args.ID) var delivery struct { ID int64 - RecipientID int64 `db:"recipient_id"` + RecipientID int64 `db:"recipient_id"` Username string Email *string SenderID int64 `db:"sender_id"` @@ -42,7 +41,7 @@ func (s *service) Handle(ctx context.Context) error { SentAt *time.Time `db:"sent_at"` } - err = pgxscan.Get(ctx, s.Conn, &delivery, ` + err := pgxscan.Get(ctx, w.Conn, &delivery, ` select d.id, r.username username, @@ -56,7 +55,7 @@ from deliveries d join users r on r.id = d.recipient_id join users s on s.id = d.sender_id join messages m on m.id = d.message_id -where d.id = $1`, j.Payload["id"]) +where d.id = $1`, job.Args.ID) if err != nil { return err } @@ -72,11 +71,11 @@ where d.id = $1`, j.Payload["id"]) link := app.AbsoluteURL(url.URL{Path: fmt.Sprintf("/deliveries/%d", delivery.ID)}) if delivery.Email == nil { - recipient, err := s.Queries.UserByID(ctx, delivery.RecipientID) + recipient, err := w.Queries.UserByID(ctx, delivery.RecipientID) if err != nil { return err } - parents, err := s.Queries.ParentsByKidID(ctx, recipient.ID) + parents, err := w.Queries.ParentsByKidID(ctx, recipient.ID) if err != nil { return err } diff --git a/worker/notifyfriend/notifyfriend.go b/worker/notifyfriend/notifyfriend.go index 1f0ecdf4..eb53647b 100644 --- a/worker/notifyfriend/notifyfriend.go +++ b/worker/notifyfriend/notifyfriend.go @@ -11,27 +11,26 @@ import ( "oj/services/email" "time" - "github.com/acaloiaro/neoq/jobs" "github.com/georgysavva/scany/v2/pgxscan" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" + "github.com/riverqueue/river" ) -type service struct { - Conn *pgxpool.Pool - Queries *api.Queries +type NotifyFriendArgs struct { + ID int64 `json:"id"` } -func NewService(q *api.Queries, conn *pgxpool.Pool) *service { - return &service{Queries: q, Conn: conn} +func (NotifyFriendArgs) Kind() string { return "notify_friend" } + +type Worker struct { + river.WorkerDefaults[NotifyFriendArgs] + Queries *api.Queries + Conn *pgxpool.Pool } -func (s *service) Handle(ctx context.Context) error { - j, err := jobs.FromContext(ctx) - if err != nil { - return err - } - log.Printf("handleNotifyFriend job id: %d, payload: %v", j.ID, j.Payload) +func (w *Worker) Work(ctx context.Context, job *river.Job[NotifyFriendArgs]) error { + log.Printf("handleNotifyFriend job id: %d, friend id: %d", job.ID, job.Args.ID) var friend struct { ID int64 @@ -43,7 +42,7 @@ func (s *service) Handle(ctx context.Context) error { TargetEmail string `db:"target_email"` } - err = pgxscan.Get(ctx, s.Conn, &friend, ` + err := pgxscan.Get(ctx, w.Conn, &friend, ` select f.id, f.created_at, a.id a_id, a.email, a.username, @@ -52,13 +51,13 @@ from friends f join users a on a.id = f.a_id join users b on b.id = f.b_id where f.id = $1 -`, j.Payload["id"]) +`, job.Args.ID) if err != nil { return err } var mutualID int64 - err = pgxscan.Get(ctx, s.Conn, &mutualID, `select id from friends where a_id = $1 and b_id = $2`, friend.BID, friend.AID) + err = pgxscan.Get(ctx, w.Conn, &mutualID, `select id from friends where a_id = $1 and b_id = $2`, friend.BID, friend.AID) if err != nil && !errors.Is(err, pgx.ErrNoRows) { return err } diff --git a/worker/notifykidfriend/notifykidfriend.go b/worker/notifykidfriend/notifykidfriend.go index fb4aba2d..f6c8c5da 100644 --- a/worker/notifykidfriend/notifykidfriend.go +++ b/worker/notifykidfriend/notifykidfriend.go @@ -11,26 +11,26 @@ import ( "oj/services/email" "time" - "github.com/acaloiaro/neoq/jobs" "github.com/georgysavva/scany/v2/pgxscan" "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/riverqueue/river" ) -type service struct { - Queries *api.Queries - Conn pgxscan.Querier +type NotifyKidFriendArgs struct { + ID int64 `json:"id"` } -func NewService(q *api.Queries, conn pgxscan.Querier) *service { - return &service{Queries: q, Conn: conn} +func (NotifyKidFriendArgs) Kind() string { return "notify_kid_friend" } + +type Worker struct { + river.WorkerDefaults[NotifyKidFriendArgs] + Queries *api.Queries + Conn *pgxpool.Pool } -func (s *service) Handle(ctx context.Context) error { - j, err := jobs.FromContext(ctx) - if err != nil { - return err - } - log.Printf("handleNotifyKidFriend job id: %d, payload: %v", j.ID, j.Payload) +func (w *Worker) Work(ctx context.Context, job *river.Job[NotifyKidFriendArgs]) error { + log.Printf("handleNotifyKidFriend job id: %d, friend id: %d", job.ID, job.Args.ID) var friend struct { ID int64 @@ -41,7 +41,7 @@ func (s *service) Handle(ctx context.Context) error { BUsername string `db:"b_username"` } - err = pgxscan.Get(ctx, s.Conn, &friend, ` + err := pgxscan.Get(ctx, w.Conn, &friend, ` select f.id, f.created_at, a.id a_id, @@ -52,20 +52,20 @@ from friends f join users a on a.id = f.a_id join users b on b.id = f.b_id where f.id = $1 -`, j.Payload["id"]) +`, job.Args.ID) if err != nil { return fmt.Errorf("getting friend %w", err) } var mutualID int64 - err = pgxscan.Get(ctx, s.Conn, &mutualID, `select id from friends where a_id = $1 and b_id = $2`, friend.BID, friend.AID) + err = pgxscan.Get(ctx, w.Conn, &mutualID, `select id from friends where a_id = $1 and b_id = $2`, friend.BID, friend.AID) if err != nil && !errors.Is(err, pgx.ErrNoRows) { return fmt.Errorf("getting mutual %w", err) } aUserLink := app.AbsoluteURL(url.URL{Path: fmt.Sprintf("/u/%d", friend.AID)}) - bParents, err := s.Queries.ParentsByKidID(ctx, friend.BID) + bParents, err := w.Queries.ParentsByKidID(ctx, friend.BID) if err != nil { return fmt.Errorf("GetParents %w", err) } @@ -86,7 +86,7 @@ where f.id = $1 } } - aParents, err := s.Queries.ParentsByKidID(ctx, friend.AID) + aParents, err := w.Queries.ParentsByKidID(ctx, friend.AID) if err != nil { return err } diff --git a/worker/worker.go b/worker/worker.go index 05d1ca0e..acee3fb0 100644 --- a/worker/worker.go +++ b/worker/worker.go @@ -2,6 +2,7 @@ package worker import ( "context" + "fmt" "log" "oj/api" "oj/worker/helloworld" @@ -9,36 +10,24 @@ import ( "oj/worker/notifyfriend" "oj/worker/notifykidfriend" "oj/worker/youtubedownload" - "time" - "github.com/acaloiaro/neoq" - "github.com/acaloiaro/neoq/handler" - "github.com/acaloiaro/neoq/jobs" - "github.com/acaloiaro/neoq/types" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" "github.com/riverqueue/river" "github.com/riverqueue/river/riverdriver/riverpgxv5" ) -var Queue types.Backend var RiverClient *river.Client[pgx.Tx] func Start(ctx context.Context, queries *api.Queries, conn *pgxpool.Pool) error { - var err error - Queue, err = neoq.New(ctx) - if err != nil { - return err - } - - Queue.Start(ctx, "notify-delivery", handler.New(notifydelivery.NewService(queries, conn).Handle)) - Queue.Start(ctx, "notify-friend", handler.New(notifyfriend.NewService(queries, conn).Handle)) - Queue.Start(ctx, "notify-kid-friend", handler.New(notifykidfriend.NewService(queries, conn).Handle)) - workers := river.NewWorkers() river.AddWorker(workers, &helloworld.Worker{}) river.AddWorker(workers, youtubedownload.NewWorker(queries)) + river.AddWorker(workers, ¬ifydelivery.Worker{Queries: queries, Conn: conn}) + river.AddWorker(workers, ¬ifyfriend.Worker{Queries: queries, Conn: conn}) + river.AddWorker(workers, ¬ifykidfriend.Worker{Queries: queries, Conn: conn}) + var err error RiverClient, err = river.NewClient(riverpgxv5.New(conn), &river.Config{ Queues: map[string]river.QueueConfig{ river.QueueDefault: {MaxWorkers: 1}, @@ -59,25 +48,27 @@ func Start(ctx context.Context, queries *api.Queries, conn *pgxpool.Pool) error } func NotifyDelivery(deliveryID int64) (string, error) { - return Queue.Enqueue(context.Background(), &jobs.Job{ - Queue: "notify-delivery", - Payload: map[string]any{"id": deliveryID}, - RunAfter: time.Now().Add(1 * time.Second), - }) -} - -func NotifyFriend(friendID int64) (string, error) { - log.Printf("Enqueue NotifyFriend %d", friendID) - return Queue.Enqueue(context.Background(), &jobs.Job{ - Queue: "notify-friend", - Payload: map[string]any{"id": friendID}, - }) + job, err := RiverClient.Insert(context.Background(), ¬ifydelivery.NotifyDeliveryArgs{ID: deliveryID}, nil) + if err != nil { + return "", err + } + return fmt.Sprint(job.Job.ID), nil } func NotifyKidFriend(friendID int64) (string, error) { log.Printf("Enqueue NotifyKidFriend %d", friendID) - return Queue.Enqueue(context.Background(), &jobs.Job{ - Queue: "notify-kid-friend", - Payload: map[string]any{"id": friendID}, - }) + job, err := RiverClient.Insert(context.Background(), ¬ifykidfriend.NotifyKidFriendArgs{ID: friendID}, nil) + if err != nil { + return "", err + } + return fmt.Sprint(job.Job.ID), nil +} + +func NotifyFriend(friendID int64) (string, error) { + log.Printf("Enqueue NotifyFriend %d", friendID) + job, err := RiverClient.Insert(context.Background(), ¬ifyfriend.NotifyFriendArgs{ID: friendID}, nil) + if err != nil { + return "", err + } + return fmt.Sprint(job.Job.ID), nil }