Commit a806da23 authored by Jacob Vosmaer (GitLab)'s avatar Jacob Vosmaer (GitLab)

Merge branch '111-watchkey-fix' into 'master'

Fix for internal/redis.TestWatchKey{AlreadyChanged,MassiveParallel}

Closes #111

See merge request !136
parents e798a664 951f21c8
...@@ -36,153 +36,128 @@ func createUnsubscribeMessage(key string) []interface{} { ...@@ -36,153 +36,128 @@ func createUnsubscribeMessage(key string) []interface{} {
} }
} }
func countWatchers(key string) int {
keyWatcherMutex.Lock()
defer keyWatcherMutex.Unlock()
return len(keyWatcher[key])
}
func deleteWatchers(key string) {
keyWatcherMutex.Lock()
defer keyWatcherMutex.Unlock()
delete(keyWatcher, key)
}
// Forces a run of the `Process` loop against a mock PubSubConn.
func processMessages(numWatchers int, value string) {
psc := redigomock.NewConn()
// Setup the initial subscription message
psc.Command("SUBSCRIBE", keySubChannel).Expect(createSubscribeMessage(keySubChannel))
psc.Command("UNSUBSCRIBE", keySubChannel).Expect(createUnsubscribeMessage(keySubChannel))
psc.AddSubscriptionMessage(createSubscriptionMessage(keySubChannel, runnerKey+"="+value))
// Wait for all the `WatchKey` calls to be registered
for countWatchers(runnerKey) != numWatchers {
time.Sleep(time.Millisecond)
}
processInner(psc)
}
func TestWatchKeySeenChange(t *testing.T) { func TestWatchKeySeenChange(t *testing.T) {
mconn, td := setupMockPool() conn, td := setupMockPool()
defer td() defer td()
go Process(false) conn.Command("GET", runnerKey).Expect("something")
// Setup the initial subscription message
mconn.Command("SUBSCRIBE", keySubChannel). wg := &sync.WaitGroup{}
Expect(createSubscribeMessage(keySubChannel)) wg.Add(1)
mconn.Command("UNSUBSCRIBE", keySubChannel).
Expect(createUnsubscribeMessage(keySubChannel)) go func() {
mconn.Command("GET", runnerKey). val, err := WatchKey(runnerKey, "something", time.Second)
Expect("something"). assert.NoError(t, err, "Expected no error")
Expect("somethingelse") assert.Equal(t, WatchKeyStatusSeenChange, val, "Expected value to change")
mconn.ReceiveWait = true wg.Done()
}()
mconn.AddSubscriptionMessage(createSubscriptionMessage(keySubChannel, runnerKey+"=somethingelse"))
processMessages(1, "somethingelse")
// ACTUALLY Fill the buffers wg.Wait()
go func(mconn *redigomock.Conn) {
mconn.ReceiveNow <- true
mconn.ReceiveNow <- true
mconn.ReceiveNow <- true
}(mconn)
val, err := WatchKey(runnerKey, "something", time.Duration(1*time.Second))
assert.NoError(t, err, "Expected no error")
assert.Equal(t, WatchKeyStatusSeenChange, val, "Expected value to change")
} }
func TestWatchKeyNoChange(t *testing.T) { func TestWatchKeyNoChange(t *testing.T) {
mconn, td := setupMockPool() conn, td := setupMockPool()
defer td() defer td()
go Process(false) conn.Command("GET", runnerKey).Expect("something")
// Setup the initial subscription message
mconn.Command("SUBSCRIBE", keySubChannel). wg := &sync.WaitGroup{}
Expect(createSubscribeMessage(keySubChannel)) wg.Add(1)
mconn.Command("UNSUBSCRIBE", keySubChannel).
Expect(createUnsubscribeMessage(keySubChannel)) go func() {
mconn.Command("GET", runnerKey). val, err := WatchKey(runnerKey, "something", time.Second)
Expect("something"). assert.NoError(t, err, "Expected no error")
Expect("something") assert.Equal(t, WatchKeyStatusNoChange, val, "Expected notification without change to value")
mconn.ReceiveWait = true wg.Done()
}()
mconn.AddSubscriptionMessage(createSubscriptionMessage(keySubChannel, runnerKey+"=something"))
processMessages(1, "something")
// ACTUALLY Fill the buffers wg.Wait()
go func(mconn *redigomock.Conn) {
mconn.ReceiveNow <- true
mconn.ReceiveNow <- true
mconn.ReceiveNow <- true
}(mconn)
val, err := WatchKey(runnerKey, "something", time.Duration(1*time.Second))
assert.NoError(t, err, "Expected no error")
assert.Equal(t, WatchKeyStatusNoChange, val, "Expected notification without change to value")
} }
func TestWatchKeyTimeout(t *testing.T) { func TestWatchKeyTimeout(t *testing.T) {
mconn, td := setupMockPool() conn, td := setupMockPool()
defer td() defer td()
go Process(false) conn.Command("GET", runnerKey).Expect("something")
// Setup the initial subscription message
mconn.Command("SUBSCRIBE", keySubChannel). val, err := WatchKey(runnerKey, "something", time.Millisecond)
Expect(createSubscribeMessage(keySubChannel))
mconn.Command("UNSUBSCRIBE", keySubChannel).
Expect(createUnsubscribeMessage(keySubChannel))
mconn.Command("GET", runnerKey).
Expect("something").
Expect("something")
mconn.ReceiveWait = true
// ACTUALLY Fill the buffers
go func(mconn *redigomock.Conn) {
mconn.ReceiveNow <- true
mconn.ReceiveNow <- true
mconn.ReceiveNow <- true
}(mconn)
val, err := WatchKey(runnerKey, "something", time.Duration(1*time.Second))
assert.NoError(t, err, "Expected no error") assert.NoError(t, err, "Expected no error")
assert.Equal(t, WatchKeyStatusTimeout, val, "Expected value to not change") assert.Equal(t, WatchKeyStatusTimeout, val, "Expected value to not change")
// Clean up watchers since Process isn't doing that for us (not running)
deleteWatchers(runnerKey)
} }
func TestWatchKeyAlreadyChanged(t *testing.T) { func TestWatchKeyAlreadyChanged(t *testing.T) {
mconn, td := setupMockPool() conn, td := setupMockPool()
defer td() defer td()
go Process(false) conn.Command("GET", runnerKey).Expect("somethingelse")
// Setup the initial subscription message
mconn.Command("SUBSCRIBE", keySubChannel). val, err := WatchKey(runnerKey, "something", time.Second)
Expect(createSubscribeMessage(keySubChannel))
mconn.Command("UNSUBSCRIBE", keySubChannel).
Expect(createUnsubscribeMessage(keySubChannel))
mconn.Command("GET", runnerKey).
Expect("somethingelse").
Expect("somethingelse")
mconn.ReceiveWait = true
// ACTUALLY Fill the buffers
go func(mconn *redigomock.Conn) {
mconn.ReceiveNow <- true
mconn.ReceiveNow <- true
mconn.ReceiveNow <- true
}(mconn)
val, err := WatchKey(runnerKey, "something", time.Duration(1*time.Second))
assert.NoError(t, err, "Expected no error") assert.NoError(t, err, "Expected no error")
assert.Equal(t, WatchKeyStatusAlreadyChanged, val, "Expected value to have already changed") assert.Equal(t, WatchKeyStatusAlreadyChanged, val, "Expected value to have already changed")
// Clean up watchers since Process isn't doing that for us (not running)
deleteWatchers(runnerKey)
} }
func TestWatchKeyMassiveParallel(t *testing.T) { func TestWatchKeyMassivelyParallel(t *testing.T) {
mconn, td := setupMockPool() runTimes := 100 // 100 parallel watchers
conn, td := setupMockPool()
defer td() defer td()
go Process(false) wg := &sync.WaitGroup{}
// Setup the initial subscription message wg.Add(runTimes)
mconn.Command("SUBSCRIBE", keySubChannel).
Expect(createSubscribeMessage(keySubChannel)) getCmd := conn.Command("GET", runnerKey)
mconn.Command("UNSUBSCRIBE", keySubChannel).
Expect(createUnsubscribeMessage(keySubChannel))
getCmd := mconn.Command("GET", runnerKey)
mconn.ReceiveWait = true
const runTimes = 100
for i := 0; i < runTimes; i++ { for i := 0; i < runTimes; i++ {
mconn.AddSubscriptionMessage(createSubscriptionMessage(keySubChannel, runnerKey+"=somethingelse"))
getCmd = getCmd.Expect("something") getCmd = getCmd.Expect("something")
} }
wg := &sync.WaitGroup{}
// Race-conditions /o/ \o\
for i := 0; i < runTimes; i++ { for i := 0; i < runTimes; i++ {
wg.Add(1) go func() {
go func(mconn *redigomock.Conn) { val, err := WatchKey(runnerKey, "something", time.Second)
defer wg.Done()
// ACTUALLY Fill the buffers
go func(mconn *redigomock.Conn) {
mconn.ReceiveNow <- true
}(mconn)
val, err := WatchKey(runnerKey, "something", time.Duration(1*time.Second))
assert.NoError(t, err, "Expected no error") assert.NoError(t, err, "Expected no error")
assert.Equal(t, WatchKeyStatusSeenChange, val, "Expected value to change") assert.Equal(t, WatchKeyStatusSeenChange, val, "Expected value to change")
}(mconn) wg.Done()
}()
} }
wg.Wait()
processMessages(runTimes, "somethingelse")
wg.Wait()
} }
...@@ -149,8 +149,7 @@ func GetString(key string) (string, error) { ...@@ -149,8 +149,7 @@ func GetString(key string) (string, error) {
if conn == nil { if conn == nil {
return "", fmt.Errorf("Not connected to redis") return "", fmt.Errorf("Not connected to redis")
} }
defer func() { defer conn.Close()
conn.Close()
}()
return redis.String(conn.Do("GET", key)) return redis.String(conn.Do("GET", key))
} }
Markdown is supported
0%
or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment