packages feed

nri-kafka-0.4.0.1: scripts/pause-resume-bug/Consumer.hs

module Consumer where

import Control.Concurrent (forkIO, threadDelay)
import Control.Concurrent.MVar (MVar, newEmptyMVar, newMVar, putMVar, tryTakeMVar, withMVar)
import Control.Monad (void)
import qualified Environment
import qualified Kafka.Worker as Kafka
import Message
import System.Environment (setEnv)
import System.IO (Handle, hPutStrLn, stderr, stdout)
import Prelude (IO, String, show)

main :: IO ()
main = do
  -- Disable console logging to make it easier to spot the bug
  setEnv "LOG_ENABLED_LOGGERS" "file"
  setEnv "LOG_FILE" "/dev/null"

  -- Reduce buffer and batch sizes to make it fail faster
  setEnv "KAFKA_MAX_MSGS_PER_PARTITION_BUFFERED_LOCALLY" "20"
  setEnv "KAFKA_POLL_BATCH_SIZE" "5"

  settings <- Environment.decode Kafka.decoder
  doAnythingHandler <- Platform.doAnythingHandler
  lastId <- newEmptyMVar

  lock <- newMVar ()

  let processMsg (msg :: Message) =
        ( do
            let msgId = ("ID(" ++ show (id msg) ++ ")")
            prevId <- tryTakeMVar lastId

            case (prevId, id msg) of
              (Nothing, _) ->
                printAtomic lock stdout (msgId ++ " First message has been received")
              (_, 1) ->
                printAtomic lock stdout (msgId ++ " Producer has been restarted")
              (Just prev, curr)
                | prev + 1 == curr ->
                    -- This is the expected behavior
                    printAtomic lock stdout (msgId ++ " OK")
              (Just prev, curr) ->
                -- This is the bug
                printAtomic
                  lock
                  stderr
                  ( "ERROR: Expected ID "
                      ++ show (prev + 1)
                      ++ " but got "
                      ++ show curr
                  )

            putMVar lastId (id msg)
            threadDelay 200000
        )
          |> fmap Ok
          |> Platform.doAnything doAnythingHandler
  let subscription = Kafka.subscription "pause-resume-bug" processMsg

  Kafka.process settings "pause-resume-bug-consumer" subscription

printAtomic :: MVar () -> Handle -> String -> IO ()
printAtomic lock handle msg = do
  (\_ -> hPutStrLn handle msg)
    |> withMVar lock
    |> forkIO
    |> void