TestService_WithDrops and TestService_SubmitVerificationWithDrops submitted three items into a size-1 queue and asserted at least one was dropped, relying on the single consumer not draining the queue between submits. synctest does not fully pin this because the MockDest send path does real logging I/O outside the bubble, so under CI's -race scheduling the consumer occasionally drained all three, delivering everything and failing the "<= 2" assertion (~0.2% of CI runs). Add an optional gate channel to MockDest so a destination blocks in Send/SendVerification until released. The tests now submit one item (the consumer picks it up and blocks on the gate), fill the size-1 queue, submit an overflow item that is dropped, then release the gate — the drop is deterministic regardless of scheduling. Assert exactly two delivered instead of "at most two".
85 lines
1.8 KiB
Go
85 lines
1.8 KiB
Go
package notify
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"sync"
|
|
|
|
log "github.com/go-pkgz/lgr"
|
|
)
|
|
|
|
// MockDest is a destination mock
|
|
type MockDest struct {
|
|
data []Request
|
|
verificationData []VerificationRequest
|
|
id int
|
|
closed bool
|
|
lock sync.Mutex
|
|
block chan struct{} // if non-nil, Send/SendVerification wait on it before recording, letting tests pin the consumer
|
|
}
|
|
|
|
// Send mock
|
|
func (m *MockDest) Send(ctx context.Context, r Request) error {
|
|
if m.block != nil {
|
|
<-m.block
|
|
}
|
|
m.lock.Lock()
|
|
defer m.lock.Unlock()
|
|
if err := ctx.Err(); err != nil {
|
|
log.Printf("ctx closed %d", m.id)
|
|
m.closed = true
|
|
return nil
|
|
}
|
|
m.data = append(m.data, r)
|
|
log.Printf("sent %s -> %d", r.Comment.ID, m.id)
|
|
return nil
|
|
}
|
|
|
|
// SendVerification mock
|
|
func (m *MockDest) SendVerification(ctx context.Context, v VerificationRequest) error {
|
|
if m.block != nil {
|
|
<-m.block
|
|
}
|
|
m.lock.Lock()
|
|
defer m.lock.Unlock()
|
|
if err := ctx.Err(); err != nil {
|
|
log.Printf("verification ctx closed %d", m.id)
|
|
m.closed = true
|
|
return nil
|
|
}
|
|
m.verificationData = append(m.verificationData, v)
|
|
log.Printf("sent verification %s -> %d", v.User, m.id)
|
|
return nil
|
|
}
|
|
|
|
// Get mock
|
|
func (m *MockDest) Get() []Request {
|
|
m.lock.Lock()
|
|
defer m.lock.Unlock()
|
|
res := make([]Request, len(m.data))
|
|
copy(res, m.data)
|
|
return res
|
|
}
|
|
|
|
// GetVerify mock
|
|
func (m *MockDest) GetVerify() []VerificationRequest {
|
|
m.lock.Lock()
|
|
defer m.lock.Unlock()
|
|
res := make([]VerificationRequest, len(m.verificationData))
|
|
copy(res, m.verificationData)
|
|
return res
|
|
}
|
|
|
|
// IsClosed returns closed status safely
|
|
func (m *MockDest) IsClosed() bool {
|
|
m.lock.Lock()
|
|
defer m.lock.Unlock()
|
|
return m.closed
|
|
}
|
|
|
|
func (m *MockDest) String() string {
|
|
m.lock.Lock()
|
|
defer m.lock.Unlock()
|
|
return fmt.Sprintf("mock id=%d, closed=%v", m.id, m.closed)
|
|
}
|