diff --git a/examples/cmd/markdownRenderer/store/queue.go b/examples/cmd/markdownRenderer/store/queue.go
index be13e411..e9c1bedf 100644
--- a/examples/cmd/markdownRenderer/store/queue.go
+++ b/examples/cmd/markdownRenderer/store/queue.go
@@ -224,35 +224,6 @@ func (conn *store) QueueGroupComplete(id int64) (done bool, cancelled bool, err
return count == 0, cancelled, err
}
-func (conn *store) IsQueueAddressComplete(address string) (done bool, err error) {
- address = strings.TrimSpace(address)
- if address == "" {
- return false, errors.New("no address provided for IsQueueAddressComplete")
- }
-
- // Include the address in the transaction description to avoid races when
- // testing with a special locking connection.
- var count int64
- var workErr error
- conn.db.Transaction(func(tx *gorm.DB) error {
- err = tx.Model(&Queue{}).Where("address = ?", address).Count(&count).Error
- if err != nil {
- return err
- }
-
- c2 := &store{
- db: tx,
- }
- workErr = c2.QueueAddressedCheck(address)
- return nil
- })
-
- if workErr != nil {
- err = workErr
- }
- return count == 0, err
-}
-
func (conn *store) IsQueueAddressInProgress(address string) (bool, error) {
address = strings.TrimSpace(address)
if address == "" {
@@ -568,87 +539,6 @@ func (conn *store) QueueDelete(permitId permit.Permit) (err error) {
})
}
-func (conn *store) QueueAddressedComplete(address string, failure error) (err error) {
-
- var bytes []byte
- if failure != nil {
- queueError, isQueueError := failure.(*queue.QueueError)
- if isQueueError {
- // Record type of queue.QueueError
- bytes, err = json.Marshal(queueError)
- } else {
- // Record generic error
- bytes, err = json.Marshal(&queue.QueueError{
- Message: failure.Error(),
- })
- }
- }
- // Return if there were any marshaling errors
- if err != nil {
- return err
- }
-
- var postgres bool
-
- err = conn.db.Transaction(func(tx *gorm.DB) error {
- if postgres {
- // Lock to prevent race. Don't allow reading until this operation is complete.
- err = tx.Exec("LOCK queue_failure IN ACCESS EXCLUSIVE MODE").Error
- if err != nil {
- return err
- }
- }
-
- // First, delete any references to this address
- err = tx.Delete(&QueueFailure{}, "address = ?", address).Error
- if err != nil {
- return err
- }
-
- // Next, if there is an error to report, insert it
- if failure != nil {
- failure := QueueFailure{
- Address: address,
- Error: string(bytes),
- }
- err = tx.Create(&failure).Error
- if err != nil {
- return err
- }
- }
- return nil
- })
- return
-}
-
-type QueueAddressFailure struct {
- Message string `json:"error"`
-}
-
-func (err *QueueAddressFailure) Error() string {
- return err.Message
-}
-
-func (conn *store) QueueAddressedCheck(address string) error {
- var failure QueueFailure
- err := conn.db.First(&failure, "address = ?", address).Error
- if err == gorm.ErrRecordNotFound {
- return nil
- } else if err != nil {
- return err
- } else {
- if failure.Error != "" {
- var queueError queue.QueueError
- if err := json.Unmarshal([]byte(failure.Error), &queueError); err != nil {
- return fmt.Errorf("error unmarshalling queue.QueueError: %s", err)
- } else {
- return &queueError
- }
- }
- }
- return nil
-}
-
func (conn *store) QueuePermits(name string) ([]dbqueuetypes.QueuePermit, error) {
permits := make([]QueuePermit, 0)
diff --git a/graph.svg b/graph.svg
index 9d417f72..b1f41b46 100644
--- a/graph.svg
+++ b/graph.svg
@@ -84,7 +84,7 @@
rselection/impls/pgx
rselection/impls/pgx
-v0.1.2
+v0.1.3
@@ -98,7 +98,7 @@
rsnotify/listeners/postgrespgx
rsnotify/listeners/postgrespgx
-v1.5.1
+v1.5.3
@@ -126,7 +126,7 @@
rslog
rslog
-v1.3.0
+v1.6.0
@@ -140,7 +140,7 @@
rsnotify/listeners/postgrespgx->rsnotify
-v1.5.1
+v1.5.2
@@ -182,7 +182,7 @@
rsqueue
rsqueue
-v0.2.1
+v0.3.0
@@ -196,14 +196,14 @@
rsqueue->rsnotify
-v1.4.0
+v1.5.1
rsstorage/servers/file
rsstorage/servers/file
-v0.3.1
+v0.3.3
diff --git a/pkg/rsqueue/impls/database/dbqueuetypes/types.go b/pkg/rsqueue/impls/database/dbqueuetypes/types.go
index 4a3d32f4..a05031b1 100644
--- a/pkg/rsqueue/impls/database/dbqueuetypes/types.go
+++ b/pkg/rsqueue/impls/database/dbqueuetypes/types.go
@@ -71,14 +71,6 @@ type QueueStore interface {
// * error - errors
QueuePeek(types ...uint64) (results []queue.QueueWork, err error)
- // IsQueueAddressComplete checks to see if an address is done/gone
- // Expects:
- // * id - the queue item address
- // Returns:
- // * bool - is the item done/gone?
- // * error - errors
- IsQueueAddressComplete(address string) (bool, error)
-
// IsQueueAddressInProgress checks to see if an address is still in progress
// Expects:
// * id - the queue item address
@@ -86,9 +78,6 @@ type QueueStore interface {
// * bool - is the item still in progress?
// * error - errors
IsQueueAddressInProgress(address string) (bool, error)
-
- // QueueAddressedComplete saves or clears failure information for an addressed item
- QueueAddressedComplete(address string, failure error) error
}
type QueueGroupStore interface {
diff --git a/pkg/rsqueue/impls/database/go.mod b/pkg/rsqueue/impls/database/go.mod
index 50aa5de8..e9e6651f 100644
--- a/pkg/rsqueue/impls/database/go.mod
+++ b/pkg/rsqueue/impls/database/go.mod
@@ -5,9 +5,9 @@ go 1.20
require (
github.com/fortytw2/leaktest v1.3.0
github.com/google/uuid v1.1.2
- github.com/rstudio/platform-lib/pkg/rsnotify v1.4.0
+ github.com/rstudio/platform-lib/pkg/rsnotify v1.5.1
github.com/rstudio/platform-lib/pkg/rsnotify/listeners/local v1.4.1
- github.com/rstudio/platform-lib/pkg/rsqueue v0.1.1
+ github.com/rstudio/platform-lib/pkg/rsqueue v0.3.0
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c
)
diff --git a/pkg/rsqueue/impls/database/go.sum b/pkg/rsqueue/impls/database/go.sum
index bf867fe5..d6212af1 100644
--- a/pkg/rsqueue/impls/database/go.sum
+++ b/pkg/rsqueue/impls/database/go.sum
@@ -9,13 +9,11 @@ github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ=
github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI=
github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
-github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLANZcx1PVRCS0qkT7pwLkGfwJo4zjcN/Tysno=
-github.com/rstudio/platform-lib/pkg/rsnotify v1.4.0 h1:tMeO7PckDZzDROBmu5nZJ1pORN3dNDWMWv57X0xpMf4=
-github.com/rstudio/platform-lib/pkg/rsnotify v1.4.0/go.mod h1:OTMZNgESF0Y2THisqq2gd4P/yF9YLCUW1GNXy5ainYM=
+github.com/rstudio/platform-lib/pkg/rsnotify v1.5.1 h1:hkLLKye0iBqW5r5ae6qnX5viU9qj0Eiwh+vUpkuy2zI=
+github.com/rstudio/platform-lib/pkg/rsnotify v1.5.1/go.mod h1:OTMZNgESF0Y2THisqq2gd4P/yF9YLCUW1GNXy5ainYM=
github.com/rstudio/platform-lib/pkg/rsnotify/listeners/local v1.4.1 h1:65BsZFFHKW9Bm0MBi3Jw7a72C2XqRnmA+519Cm+6zFk=
github.com/rstudio/platform-lib/pkg/rsnotify/listeners/local v1.4.1/go.mod h1:3fj43uPHSrY30ZJmAgbSmmy9d7heQVnL/MXVwOz+5Pc=
-github.com/rstudio/platform-lib/pkg/rsqueue v0.1.1 h1:y6hxq3z1F98KFEnfSSYDm5BP4C9MlSxwnNPKHGluV7w=
-github.com/rstudio/platform-lib/pkg/rsqueue v0.1.1/go.mod h1:v47Q+KI/vBl6NLgrMLry6x3J651uhppVvCCGffq/bDI=
-gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
+github.com/rstudio/platform-lib/pkg/rsqueue v0.3.0 h1:2/risfO4VatW7FcmWeup0Mbx1QlPLYYUZrN8sb8Pt+g=
+github.com/rstudio/platform-lib/pkg/rsqueue v0.3.0/go.mod h1:ac7W+fPG4pcb8yYp/npchGewPJezAu6UjlN0DjPEuH8=
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=
diff --git a/pkg/rsqueue/impls/database/groupqueue_test.go b/pkg/rsqueue/impls/database/groupqueue_test.go
index b79bbd72..7bcdb46b 100644
--- a/pkg/rsqueue/impls/database/groupqueue_test.go
+++ b/pkg/rsqueue/impls/database/groupqueue_test.go
@@ -103,9 +103,6 @@ func (f *fakeQueue) IsAddressInQueue(address string) (bool, error) {
func (f *fakeQueue) PollAddress(address string) (errs <-chan error) {
return f.pollErrs
}
-func (f *fakeQueue) RecordFailure(address string, failure error) error {
- return f.record
-}
func (f *fakeQueue) Get(maxPriority uint64, maxPriorityChan chan uint64, types queue.QueueSupportedTypes, stop chan bool) (*queue.QueueWork, error) {
return nil, errors.New("n/i")
}
diff --git a/pkg/rsqueue/impls/database/queue.go b/pkg/rsqueue/impls/database/queue.go
index 85163622..5cdb700a 100644
--- a/pkg/rsqueue/impls/database/queue.go
+++ b/pkg/rsqueue/impls/database/queue.go
@@ -234,10 +234,6 @@ func (q *DatabaseQueue) AddressedPush(priority uint64, groupId int64, address st
return err
}
-func (q *DatabaseQueue) RecordFailure(address string, failure error) error {
- return q.store.QueueAddressedComplete(address, failure)
-}
-
func (q *DatabaseQueue) IsAddressInQueue(address string) (bool, error) {
return q.store.IsQueueAddressInProgress(address)
}
@@ -247,7 +243,7 @@ func (q *DatabaseQueue) PollAddress(address string) <-chan error {
go func() {
for {
- isDone, err := q.store.IsQueueAddressComplete(address)
+ queued, err := q.IsAddressInQueue(address)
if err != nil {
// Ignore lock errors
if !utils.IsSqliteLockError(err) {
@@ -255,14 +251,14 @@ func (q *DatabaseQueue) PollAddress(address string) <-chan error {
close(errCh)
return
}
- } else if isDone {
+ } else if !queued {
q.debugLogger.Debugf("Queue work with address %s completed", address)
close(errCh)
return
}
// Wait for a notification or an interval, then poll again
- done := func() bool {
+ done, err := func() (bool, error) {
completedMsgs := q.SubscribeOne(q.notifyTypeWorkComplete, func(n listener.Notification) bool {
if wn, ok := n.(*agenttypes.WorkCompleteNotification); ok {
return wn.Address == address
@@ -279,20 +275,25 @@ func (q *DatabaseQueue) PollAddress(address string) <-chan error {
defer tick.Stop()
for {
select {
- case <-completedMsgs:
+ case n := <-completedMsgs:
q.debugLogger.Debugf("Queue was notified that work with address %s completed", address)
- return false
+ if wn, ok := n.(*agenttypes.WorkCompleteNotification); ok {
+ return true, wn.Error
+ }
case <-chunkMsgs:
q.debugLogger.Debugf("Queue was notified that chunk with address %s is ready", address)
- return true
+ return true, nil
case <-tick.C:
- return false
+ return false, nil
}
}
}()
- // If we received a chunk notification, then we return immediately so the client can
- // begin downloading chunks
+ // If we received a chunk or work complete notification, then we return immediately so the client can
+ // continue.
if done {
+ if err != nil {
+ errCh <- err
+ }
close(errCh)
return
}
@@ -372,7 +373,7 @@ func (q *DatabaseQueue) Get(maxPriority uint64, maxPriorityChan chan uint64, typ
if err != nil {
return queueWork, err
}
- queueWork, err := q.store.QueuePop(q.name, maxPriority, types.Enabled())
+ queueWork, err = q.store.QueuePop(q.name, maxPriority, types.Enabled())
if err != sql.ErrNoRows {
q.measureDequeue(queueWork, err)
return queueWork, err
diff --git a/pkg/rsqueue/impls/database/queue_test.go b/pkg/rsqueue/impls/database/queue_test.go
index a6127617..0042a865 100644
--- a/pkg/rsqueue/impls/database/queue_test.go
+++ b/pkg/rsqueue/impls/database/queue_test.go
@@ -43,8 +43,6 @@ type QueueTestStore struct {
llf *local.ListenerProvider
err error
polled int
- poll bool
- pollErr error
hasAddress bool
hasAddressErr error
enabled []uint64
@@ -101,12 +99,8 @@ func (s *QueueTestStore) QueuePop(name string, maxPriority uint64, types []uint6
}
func (s *QueueTestStore) IsQueueAddressInProgress(address string) (bool, error) {
- return s.hasAddress, s.hasAddressErr
-}
-
-func (s *QueueTestStore) IsQueueAddressComplete(address string) (bool, error) {
s.polled++
- return s.poll, s.pollErr
+ return s.hasAddress, s.hasAddressErr
}
func (s *QueueTestStore) NotifyExtend(permit uint64) error {
@@ -117,10 +111,6 @@ func (s *QueueTestStore) QueueDelete(permit permit.Permit) error {
return s.err
}
-func (s *QueueTestStore) QueueAddressedComplete(address string, failure error) error {
- return s.err
-}
-
func (s *QueueTestStore) GetLocalListenerProvider() *local.ListenerProvider {
if s.llf != nil {
return s.llf
@@ -211,29 +201,6 @@ func (s *QueueSuite) TestNewQueue(c *check.C) {
_ = typeQueue
}
-func (s *QueueSuite) TestRecord(c *check.C) {
- q := &DatabaseQueue{
- store: s.store,
- debugLogger: &fakeLogger{},
- wrapper: &fakeWrapper{},
- }
- err := q.RecordFailure("abc", errors.New("test"))
- c.Assert(err, check.IsNil)
- err = q.RecordFailure("abc", nil)
- c.Assert(err, check.IsNil)
-}
-
-func (s *QueueSuite) TestRecordErrs(c *check.C) {
- s.store.err = errors.New("kaboom!")
- q := &DatabaseQueue{
- store: s.store,
- debugLogger: &fakeLogger{},
- wrapper: &fakeWrapper{},
- }
- err := q.RecordFailure("abc", errors.New("test"))
- c.Assert(err, check.NotNil)
-}
-
func (s *QueueSuite) TestPush(c *check.C) {
q := &DatabaseQueue{
store: s.store,
@@ -452,7 +419,7 @@ func (s *QueueSuite) TestDeleteErrs(c *check.C) {
func (s *QueueSuite) TestPollErr(c *check.C) {
defer leaktest.Check(c)
- s.store.pollErr = errors.New("horrible error")
+ s.store.hasAddressErr = errors.New("horrible error")
q := &DatabaseQueue{
store: s.store,
debugLogger: &fakeLogger{},
@@ -485,7 +452,7 @@ func (s *QueueSuite) TestPollErr(c *check.C) {
func (s *QueueSuite) TestPollLockErr(c *check.C) {
defer leaktest.Check(c)
- s.store.pollErr = errors.New("database is locked")
+ s.store.hasAddressErr = errors.New("database is locked")
q := &DatabaseQueue{
store: s.store,
debugLogger: &fakeLogger{},
@@ -510,8 +477,8 @@ func (s *QueueSuite) TestPollLockErr(c *check.C) {
go func() {
time.Sleep(time.Millisecond * 30)
- s.store.pollErr = nil
- s.store.poll = true
+ s.store.hasAddressErr = nil
+ s.store.hasAddress = false
}()
err := <-errCh
@@ -522,7 +489,7 @@ func (s *QueueSuite) TestPollLockErr(c *check.C) {
func (s *QueueSuite) TestPollTickOk(c *check.C) {
defer leaktest.Check(c)
- s.store.poll = false
+ s.store.hasAddress = true
q := &DatabaseQueue{
store: s.store,
debugLogger: &fakeLogger{},
@@ -547,7 +514,7 @@ func (s *QueueSuite) TestPollTickOk(c *check.C) {
go func() {
time.Sleep(time.Millisecond * 30)
- s.store.poll = true
+ s.store.hasAddress = false
}()
<-errCh
@@ -557,7 +524,9 @@ func (s *QueueSuite) TestPollTickOk(c *check.C) {
func (s *QueueSuite) TestPollNotifyOk(c *check.C) {
defer leaktest.Check(c)
- cstore := &QueueTestStore{}
+ cstore := &QueueTestStore{
+ hasAddress: true,
+ }
q := &DatabaseQueue{
store: cstore,
debugLogger: &fakeLogger{},
@@ -588,18 +557,59 @@ func (s *QueueSuite) TestPollNotifyOk(c *check.C) {
time.Sleep(10 * time.Millisecond)
// This should not initiate a poll, since the address doesn't match
- workMsgs <- agenttypes.NewWorkCompleteNotification("nothing", 10)
- // This should initiate a poll
- workMsgs <- agenttypes.NewWorkCompleteNotification("something", 10)
+ workMsgs <- agenttypes.NewWorkCompleteNotification("nothing", 10, nil)
+ // This message indicates the work is done
+ workMsgs <- agenttypes.NewWorkCompleteNotification("something", 10, nil)
- // Now the work is done
- // Sleep a bit to prevent a race when setting store.poll = true
- time.Sleep(100 * time.Millisecond)
- cstore.poll = true
- workMsgs <- agenttypes.NewWorkCompleteNotification("something", 10)
+ err := <-errCh
+ c.Assert(err, check.IsNil)
+ c.Assert(cstore.polled, check.Equals, 1)
+}
- <-errCh
- c.Assert(cstore.polled, check.Equals, 3)
+func (s *QueueSuite) TestPollNotifyErr(c *check.C) {
+ defer leaktest.Check(c)
+
+ cstore := &QueueTestStore{
+ hasAddress: true,
+ }
+ q := &DatabaseQueue{
+ store: cstore,
+ debugLogger: &fakeLogger{},
+ addressPollInterval: time.Hour,
+ subscribe: make(chan broadcaster.Subscription),
+ unsubscribe: make(chan (<-chan listener.Notification)),
+ wrapper: &fakeWrapper{},
+
+ notifyTypeWorkReady: 9,
+ notifyTypeWorkComplete: 10,
+ notifyTypeChunk: 11,
+ chunkMatcher: &fakeMatcher{},
+ }
+
+ queueMsgs := make(chan listener.Notification)
+ workMsgs := make(chan listener.Notification)
+ chunkMsgs := make(chan listener.Notification)
+ defer close(queueMsgs)
+ defer close(workMsgs)
+ defer close(chunkMsgs)
+
+ stopper := make(chan bool)
+ defer func() { stopper <- true }()
+ go q.broadcast(stopper, queueMsgs, workMsgs, chunkMsgs)
+
+ errCh := q.PollAddress("something")
+
+ time.Sleep(10 * time.Millisecond)
+
+ // This should not initiate a poll, since the address doesn't match
+ workMsgs <- agenttypes.NewWorkCompleteNotification("nothing", 10, nil)
+ // This message indicates the work is done
+ errWork := errors.New("work error")
+ workMsgs <- agenttypes.NewWorkCompleteNotification("something", 10, errWork)
+
+ err := <-errCh
+ c.Assert(err, check.ErrorMatches, "work error")
+ c.Assert(cstore.polled, check.Equals, 1)
}
func (s *QueueSuite) TestPollNotifyChunkOk(c *check.C) {
diff --git a/pkg/rsqueue/impls/database/tasks/monitor_test.go b/pkg/rsqueue/impls/database/tasks/monitor_test.go
index da4624d7..44969d7a 100644
--- a/pkg/rsqueue/impls/database/tasks/monitor_test.go
+++ b/pkg/rsqueue/impls/database/tasks/monitor_test.go
@@ -111,10 +111,6 @@ func (s *QueueTestStore) QueueDelete(permit permit.Permit) error {
return s.err
}
-func (s *QueueTestStore) QueueAddressedComplete(address string, failure error) error {
- return s.err
-}
-
func (s *QueueTestStore) GetLocalListenerProvider() *local.ListenerProvider {
if s.llf != nil {
return s.llf