diff --git a/CHANGELOG.md b/CHANGELOG.md
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -1,3 +1,30 @@
+# 0.4.0.1
+
+- Require `nri-prelude >= 0.7.0.0`
+
+# 0.4.0.0
+
+- Support GHC 9.10.2, GHC 9.12.2, `containers-0.7.x`
+- Drop support for GHC 9.4.x
+
+# 0.3.0.0
+
+- Drop support for GHC 9.2.x
+- Support GHC 9.8.3, `bytestring-0.12.x.x`, `text-2.1.x`, `aeson-2.2.x.x`
+
+# 0.2.0.1
+
+- Drop support for `aeson-1.x`
+- Support GHC 9.6.5
+
+# 0.2.0.0
+
+- Drop support for GHC 8.10.7
+
+# 0.1.0.5
+
+- Support GHC 9.4.7, `aeson-2.1.x`
+
 # 0.1.0.4
 
 - Added new `ElsewhereButToKafkaAsWell` mode to `CommitOffsets`, which commits offsets to Kafka once the external Offset storage has been updated. Kafka commits are performed only to keep Kafka informed about consumer lag.
diff --git a/LICENSE b/LICENSE
--- a/LICENSE
+++ b/LICENSE
@@ -1,6 +1,6 @@
 BSD 3-Clause License
 
-Copyright (c) 2022, NoRedInk
+Copyright (c) 2026, NoRedInk
 All rights reserved.
 
 Redistribution and use in source and binary forms, with or without
diff --git a/nri-kafka.cabal b/nri-kafka.cabal
--- a/nri-kafka.cabal
+++ b/nri-kafka.cabal
@@ -1,11 +1,11 @@
 cabal-version: 1.12
 
--- This file has been generated from package.yaml by hpack version 0.34.5.
+-- This file has been generated from package.yaml by hpack version 0.38.3.
 --
 -- see: https://github.com/sol/hpack
 
 name:           nri-kafka
-version:        0.1.0.4
+version:        0.4.0.1
 synopsis:       Functions for working with Kafka
 description:    Please see the README at <https://github.com/NoRedInk/haskell-libraries/tree/trunk/nri-kafka#readme>.
 category:       Web
@@ -13,7 +13,7 @@
 bug-reports:    https://github.com/NoRedInk/haskell-libraries/issues
 author:         NoRedInk
 maintainer:     haskell-open-source@noredink.com
-copyright:      2022 NoRedInk Corp.
+copyright:      2026 NoRedInk Corp.
 license:        BSD3
 license-file:   LICENSE
 build-type:     Simple
@@ -27,6 +27,14 @@
   location: https://github.com/NoRedInk/haskell-libraries
   subdir: nri-kafka
 
+flag pause-resume-bug
+  manual: False
+  default: False
+
+flag sync-write-benchmark
+  manual: False
+  default: False
+
 library
   exposed-modules:
       Kafka
@@ -63,24 +71,169 @@
       TypeOperators
   ghc-options: -Wall -Wcompat -Widentities -Wincomplete-record-updates -Wpartial-fields -Wredundant-constraints -Wincomplete-uni-patterns -fno-warn-type-defaults -fplugin=NriPrelude.Plugin
   build-depends:
-      aeson >=1.4.6.0 && <2.1
+      aeson >=2.0 && <2.3
     , async >=2.2.2 && <2.3
-    , base >=4.12.0.0 && <4.17
-    , bytestring >=0.10.8.2 && <0.12
+    , base >=4.18 && <4.22
+    , bytestring >=0.10.8.2 && <0.13
     , conduit >=1.3.0 && <1.4
-    , containers >=0.6.0.1 && <0.7
+    , containers >=0.6.0.1 && <0.8
     , hw-kafka-client >=4.0.3 && <5.0
-    , nri-env-parser >=0.1.0.0 && <0.2
-    , nri-observability >=0.1.1.1 && <0.2
-    , nri-prelude >=0.1.0.0 && <0.7
+    , nri-env-parser >=0.1.0.0 && <0.5
+    , nri-observability >=0.1.1.1 && <0.5
+    , nri-prelude >=0.7.0.0 && <0.8
     , safe-exceptions >=0.1.7.0 && <1.3
     , stm >=2.4 && <2.6
-    , text >=1.2.3.1 && <2.1
+    , text >=1.2.3.1 && <2.2
     , time >=1.8.0.2 && <2
-    , unix >=2.7.2.2 && <2.8.0.0
+    , unix >=2.7.2.2 && <2.9
     , uuid >=1.3.0 && <1.4
   default-language: Haskell2010
 
+executable pause-resume-bug-consumer
+  main-is: Consumer.hs
+  other-modules:
+      Message
+      Producer
+      Paths_nri_kafka
+  hs-source-dirs:
+      scripts/pause-resume-bug
+  default-extensions:
+      DataKinds
+      DeriveGeneric
+      ExtendedDefaultRules
+      FlexibleContexts
+      FlexibleInstances
+      GeneralizedNewtypeDeriving
+      MultiParamTypeClasses
+      NamedFieldPuns
+      NoImplicitPrelude
+      NumericUnderscores
+      OverloadedStrings
+      PartialTypeSignatures
+      ScopedTypeVariables
+      Strict
+      TypeOperators
+  ghc-options: -Wall -Wcompat -Widentities -Wincomplete-record-updates -Wpartial-fields -Wredundant-constraints -Wincomplete-uni-patterns -fno-warn-type-defaults -fplugin=NriPrelude.Plugin -threaded -rtsopts "-with-rtsopts=-N -T" -O2 -main-is Consumer
+  build-depends:
+      aeson >=2.0 && <2.3
+    , async >=2.2.2 && <2.3
+    , base >=4.18 && <4.22
+    , bytestring >=0.10.8.2 && <0.13
+    , conduit >=1.3.0 && <1.4
+    , containers >=0.6.0.1 && <0.8
+    , hw-kafka-client >=4.0.3 && <5.0
+    , nri-env-parser >=0.1.0.0 && <0.5
+    , nri-kafka
+    , nri-observability >=0.1.1.1 && <0.5
+    , nri-prelude >=0.7.0.0 && <0.8
+    , safe-exceptions >=0.1.7.0 && <1.3
+    , stm >=2.4 && <2.6
+    , text >=1.2.3.1 && <2.2
+    , time >=1.8.0.2 && <2
+    , unix >=2.7.2.2 && <2.9
+    , uuid >=1.3.0 && <1.4
+  default-language: Haskell2010
+  if flag(pause-resume-bug)
+    buildable: True
+  else
+    buildable: False
+
+executable pause-resume-bug-producer
+  main-is: Producer.hs
+  other-modules:
+      Consumer
+      Message
+      Paths_nri_kafka
+  hs-source-dirs:
+      scripts/pause-resume-bug
+  default-extensions:
+      DataKinds
+      DeriveGeneric
+      ExtendedDefaultRules
+      FlexibleContexts
+      FlexibleInstances
+      GeneralizedNewtypeDeriving
+      MultiParamTypeClasses
+      NamedFieldPuns
+      NoImplicitPrelude
+      NumericUnderscores
+      OverloadedStrings
+      PartialTypeSignatures
+      ScopedTypeVariables
+      Strict
+      TypeOperators
+  ghc-options: -Wall -Wcompat -Widentities -Wincomplete-record-updates -Wpartial-fields -Wredundant-constraints -Wincomplete-uni-patterns -fno-warn-type-defaults -fplugin=NriPrelude.Plugin -main-is Producer
+  build-depends:
+      aeson >=2.0 && <2.3
+    , async >=2.2.2 && <2.3
+    , base >=4.18 && <4.22
+    , bytestring >=0.10.8.2 && <0.13
+    , conduit >=1.3.0 && <1.4
+    , containers >=0.6.0.1 && <0.8
+    , hw-kafka-client >=4.0.3 && <5.0
+    , nri-env-parser >=0.1.0.0 && <0.5
+    , nri-kafka
+    , nri-observability >=0.1.1.1 && <0.5
+    , nri-prelude >=0.7.0.0 && <0.8
+    , safe-exceptions >=0.1.7.0 && <1.3
+    , stm >=2.4 && <2.6
+    , text >=1.2.3.1 && <2.2
+    , time >=1.8.0.2 && <2
+    , unix >=2.7.2.2 && <2.9
+    , uuid >=1.3.0 && <1.4
+  default-language: Haskell2010
+  if flag(pause-resume-bug)
+    buildable: True
+  else
+    buildable: False
+
+executable sync-write-benchmark
+  main-is: Main.hs
+  other-modules:
+      Paths_nri_kafka
+  hs-source-dirs:
+      scripts/sync-write-benchmark
+  default-extensions:
+      DataKinds
+      DeriveGeneric
+      ExtendedDefaultRules
+      FlexibleContexts
+      FlexibleInstances
+      GeneralizedNewtypeDeriving
+      MultiParamTypeClasses
+      NamedFieldPuns
+      NoImplicitPrelude
+      NumericUnderscores
+      OverloadedStrings
+      PartialTypeSignatures
+      ScopedTypeVariables
+      Strict
+      TypeOperators
+  ghc-options: -Wall -Wcompat -Widentities -Wincomplete-record-updates -Wpartial-fields -Wredundant-constraints -Wincomplete-uni-patterns -fno-warn-type-defaults -fplugin=NriPrelude.Plugin -threaded -rtsopts "-with-rtsopts=-N -T" -O2 -main-is Main
+  build-depends:
+      aeson >=2.0 && <2.3
+    , async >=2.2.2 && <2.3
+    , base >=4.18 && <4.22
+    , bytestring >=0.10.8.2 && <0.13
+    , conduit >=1.3.0 && <1.4
+    , containers >=0.6.0.1 && <0.8
+    , hw-kafka-client >=4.0.3 && <5.0
+    , nri-env-parser >=0.1.0.0 && <0.5
+    , nri-kafka
+    , nri-observability >=0.1.1.1 && <0.5
+    , nri-prelude >=0.7.0.0 && <0.8
+    , safe-exceptions >=0.1.7.0 && <1.3
+    , stm >=2.4 && <2.6
+    , text >=1.2.3.1 && <2.2
+    , time >=1.8.0.2 && <2
+    , unix >=2.7.2.2 && <2.9
+    , uuid >=1.3.0 && <1.4
+  default-language: Haskell2010
+  if flag(sync-write-benchmark)
+    buildable: True
+  else
+    buildable: False
+
 test-suite tests
   type: exitcode-stdio-1.0
   main-is: Main.hs
@@ -122,20 +275,20 @@
       TypeOperators
   ghc-options: -Wall -Wcompat -Widentities -Wincomplete-record-updates -Wpartial-fields -Wredundant-constraints -Wincomplete-uni-patterns -fno-warn-type-defaults -fplugin=NriPrelude.Plugin -threaded -rtsopts "-with-rtsopts=-N -T" -fno-warn-type-defaults
   build-depends:
-      aeson >=1.4.6.0 && <2.1
+      aeson >=2.0 && <2.3
     , async >=2.2.2 && <2.3
-    , base >=4.12.0.0 && <4.17
-    , bytestring >=0.10.8.2 && <0.12
+    , base >=4.18 && <4.22
+    , bytestring >=0.10.8.2 && <0.13
     , conduit >=1.3.0 && <1.4
-    , containers >=0.6.0.1 && <0.7
+    , containers >=0.6.0.1 && <0.8
     , hw-kafka-client >=4.0.3 && <5.0
-    , nri-env-parser >=0.1.0.0 && <0.2
-    , nri-observability >=0.1.1.1 && <0.2
-    , nri-prelude >=0.1.0.0 && <0.7
+    , nri-env-parser >=0.1.0.0 && <0.5
+    , nri-observability >=0.1.1.1 && <0.5
+    , nri-prelude >=0.7.0.0 && <0.8
     , safe-exceptions >=0.1.7.0 && <1.3
     , stm >=2.4 && <2.6
-    , text >=1.2.3.1 && <2.1
+    , text >=1.2.3.1 && <2.2
     , time >=1.8.0.2 && <2
-    , unix >=2.7.2.2 && <2.8.0.0
+    , unix >=2.7.2.2 && <2.9
     , uuid >=1.3.0 && <1.4
   default-language: Haskell2010
diff --git a/scripts/pause-resume-bug/Consumer.hs b/scripts/pause-resume-bug/Consumer.hs
new file mode 100644
--- /dev/null
+++ b/scripts/pause-resume-bug/Consumer.hs
@@ -0,0 +1,68 @@
+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
diff --git a/scripts/pause-resume-bug/Message.hs b/scripts/pause-resume-bug/Message.hs
new file mode 100644
--- /dev/null
+++ b/scripts/pause-resume-bug/Message.hs
@@ -0,0 +1,12 @@
+module Message where
+
+import Data.Aeson
+
+newtype Message = Message
+  { id :: Int
+  }
+  deriving (Generic)
+
+instance FromJSON Message
+
+instance ToJSON Message
diff --git a/scripts/pause-resume-bug/Producer.hs b/scripts/pause-resume-bug/Producer.hs
new file mode 100644
--- /dev/null
+++ b/scripts/pause-resume-bug/Producer.hs
@@ -0,0 +1,36 @@
+{-# LANGUAGE DisambiguateRecordFields #-}
+
+module Producer where
+
+import Conduit
+import qualified Environment
+import qualified Kafka
+import Message
+import System.Environment (setEnv)
+import Prelude (IO, error, pure)
+
+main :: IO ()
+main = do
+  setEnv "KAFKA_COMPRESSION_CODEC" "Gzip"
+  settings <- Environment.decode Kafka.decoder
+  logHandler <- Platform.silentHandler
+
+  putTextLn "Sending messages..."
+
+  Conduit.withAcquire (Kafka.handler settings) <| \handler -> do
+    [1 .. 300]
+      |> List.map
+        ( \id -> do
+            let msg =
+                  Kafka.emptyMsg "pause-resume-bug"
+                    |> Kafka.addPayload Message {id}
+                    |> Kafka.addKey "a-partition-key"
+            Kafka.sendSync handler msg
+        )
+      |> Task.sequence
+      |> Task.attempt logHandler
+      |> andThen fromResult
+
+fromResult :: Result Text a -> IO ()
+fromResult (Ok _) = pure ()
+fromResult (Err err) = error (Text.toList err)
diff --git a/scripts/sync-write-benchmark/Main.hs b/scripts/sync-write-benchmark/Main.hs
new file mode 100644
--- /dev/null
+++ b/scripts/sync-write-benchmark/Main.hs
@@ -0,0 +1,132 @@
+{-# LANGUAGE DisambiguateRecordFields #-}
+
+-- | Benchmark for `Kafka.sendSync` per-message latency.
+--
+-- Measures wall time of N sequential synchronous produces against the broker
+-- configured by 'KAFKA_BROKER_ADDRESSES'. Run with the same env on master
+-- (baseline) and on the fix branch (after) to compare the polling-induced
+-- latency floor.
+--
+-- > cabal run sync-write-benchmark -fsync-write-benchmark
+--
+-- Tunables:
+--   BENCHMARK_TOPIC          (default: nri-kafka-sync-benchmark)
+--   BENCHMARK_MESSAGE_COUNT  (default: 1000)
+--   BENCHMARK_WARMUP         (default: 50)  -- excluded from stats
+module Main where
+
+import qualified Conduit
+import Data.Aeson (FromJSON, ToJSON)
+import qualified Data.List
+import qualified Data.Word
+import qualified Environment
+import GHC.Clock (getMonotonicTimeNSec)
+import qualified Kafka
+import qualified System.Environment
+import Text.Printf (printf)
+import Prelude (IO, String, error, fromIntegral, pure, putStrLn)
+import qualified Prelude
+
+newtype Sample = Sample {idx :: Int}
+  deriving (Generic)
+
+instance FromJSON Sample
+
+instance ToJSON Sample
+
+main :: IO ()
+main = do
+  topic <- envText "BENCHMARK_TOPIC" "nri-kafka-sync-benchmark"
+  count <- envInt "BENCHMARK_MESSAGE_COUNT" 1000
+  warmup <- envInt "BENCHMARK_WARMUP" 50
+
+  settings <- Environment.decode Kafka.decoder
+  logHandler <- Platform.silentHandler
+
+  putStrLn
+    ( printf
+        "Benchmarking sync writes: warmup=%d count=%d topic=%s"
+        warmup
+        count
+        (Text.toList topic)
+    )
+
+  Conduit.withAcquire (Kafka.handler settings) <| \handler -> do
+    putStrLn "Warming up..."
+    Prelude.mapM_ (sendOne logHandler handler topic) [1 .. warmup]
+
+    putStrLn (printf "Measuring %d sends..." count)
+    samples <-
+      Prelude.traverse
+        (timeOne logHandler handler topic)
+        [warmup + 1 .. warmup + count]
+
+    putStrLn (formatStats samples)
+
+envText :: String -> Text -> IO Text
+envText name def = do
+  v <- System.Environment.lookupEnv name
+  pure
+    ( case v of
+        Prelude.Just s -> Text.fromList s
+        Prelude.Nothing -> def
+    )
+
+envInt :: String -> Int -> IO Int
+envInt name def = do
+  v <- System.Environment.lookupEnv name
+  pure
+    ( case v of
+        Prelude.Just s -> case Prelude.reads s of
+          [(i, "")] -> i
+          _ -> def
+        Prelude.Nothing -> def
+    )
+
+mkMsg :: Text -> Int -> Kafka.Msg
+mkMsg topic i =
+  Kafka.emptyMsg topic
+    |> Kafka.addPayload (Sample {idx = i})
+    -- Pin all messages to one partition so partition variability doesn't
+    -- pollute the latency distribution.
+    |> Kafka.addKey "benchmark"
+
+sendOne :: Platform.LogHandler -> Kafka.Handler -> Text -> Int -> IO ()
+sendOne logHandler handler topic i = do
+  result <- Task.attempt logHandler (Kafka.sendSync handler (mkMsg topic i))
+  case result of
+    Ok _ -> pure ()
+    Err err -> error (Text.toList err)
+
+timeOne :: Platform.LogHandler -> Kafka.Handler -> Text -> Int -> IO Data.Word.Word64
+timeOne logHandler handler topic i = do
+  let msg = mkMsg topic i
+  t0 <- getMonotonicTimeNSec
+  result <- Task.attempt logHandler (Kafka.sendSync handler msg)
+  t1 <- getMonotonicTimeNSec
+  case result of
+    Ok _ -> pure (t1 - t0)
+    Err err -> error (Text.toList err)
+
+formatStats :: [Data.Word.Word64] -> String
+formatStats samples =
+  let sorted = Data.List.sort samples
+      n = Prelude.length sorted
+      total = Prelude.sum sorted
+      avg = nsToMs (total `Prelude.div` fromIntegral n)
+      pct :: Prelude.Double -> Prelude.Double
+      pct p =
+        let i = Prelude.min (n - 1) (Prelude.floor (Prelude.fromIntegral n Prelude.* p / 100))
+         in nsToMs (sorted Prelude.!! i)
+   in printf
+        "count=%d  min=%.1fms  avg=%.1fms  p50=%.1fms  p95=%.1fms  p99=%.1fms  max=%.1fms"
+        n
+        (nsToMs (Prelude.head sorted))
+        avg
+        (pct 50)
+        (pct 95)
+        (pct 99)
+        (nsToMs (Prelude.last sorted))
+
+nsToMs :: Data.Word.Word64 -> Prelude.Double
+nsToMs ns = fromIntegral ns / 1_000_000
diff --git a/src/Kafka.hs b/src/Kafka.hs
--- a/src/Kafka.hs
+++ b/src/Kafka.hs
@@ -30,7 +30,6 @@
 where
 
 import qualified Conduit
-import qualified Control.Concurrent
 import qualified Control.Concurrent.Async as Async
 import qualified Control.Concurrent.STM as STM
 import qualified Control.Concurrent.STM.TMVar as TMVar
@@ -176,16 +175,37 @@
       |> Async.async
   Prelude.pure terminator
 
--- | We use a little trick here to poll events, by sending an empty message batch.
--- This will call the internal pollEvent function in hw-kafka-client.
+-- | Drains librdkafka's event queue (delivery reports, errors, stats) so
+-- per-message delivery callbacks fire promptly. 'Producer.flushProducer' is
+-- internally a loop over @rd_kafka_poll(rk, 100)@ that exits when the
+-- outbound queue is empty; each underlying poll blocks the OS thread up to
+-- 100ms but returns immediately when an event arrives, so latency-sensitive
+-- 'sendSync' callers wake on event rather than on a fixed clock tick. The
+-- previous implementation used a 100ms threadDelay between non-blocking
+-- drains, which added up to a 100ms (~50ms mean) wait per sync send.
+--
+-- We can't import 'Kafka.Internal.Shared.pollEvents' directly to pass our own
+-- timeout because that module is hidden in hw-kafka-client 4.x.
+-- 'flushProducer' is the closest exposed alternative; despite its name and
+-- the @closeProducer = flushProducer@ alias, it does not transition the
+-- producer to a closed state — it just flushes. It is safe to call
+-- repeatedly on a live producer.
+--
+-- Shutdown: this loop is cancelled by 'Async.race_' when the terminator
+-- TMVar is signalled. GHC can't deliver the async exception while the
+-- thread is parked in 'safe' FFI, but 'flushProducer' is not a single FFI
+-- call — it's a Haskell loop calling @rd_kafka_poll(rk, 100)@. Each poll
+-- returns within 100ms (when idle) or immediately (on event), and control
+-- briefly returns to Haskell between iterations, where any pending async
+-- exception is delivered and kills the thread. So shutdown is bounded by
+-- ~100ms even in the worst case (idle producer, no events flowing). The
+-- 'Async.race_' cleanup itself does not wait for the loser to finish dying,
+-- so we don't deadlock the release chain.
 pollEvents :: Producer.KafkaProducer -> Prelude.IO ()
 pollEvents producer = do
-  Producer.produceMessageBatch producer []
-    |> map (\_ -> ())
-  Control.Concurrent.threadDelay 100_000 {- 100ms -}
+  Producer.flushProducer producer
   pollEvents producer
 
--- |
 mkHandler :: Settings.Settings -> Producer.KafkaProducer -> Prelude.IO Internal.Handler
 mkHandler settings producer = do
   doAnything <- Platform.doAnythingHandler
@@ -216,12 +236,12 @@
     |> Platform.doAnything doAnything
 
 mkProducer :: Settings.Settings -> Prelude.IO Producer.KafkaProducer
-mkProducer Settings.Settings {Settings.brokerAddresses, Settings.deliveryTimeout, Settings.logLevel, Settings.batchNumMessages} = do
+mkProducer Settings.Settings {Settings.brokerAddresses, Settings.deliveryTimeout, Settings.logLevel, Settings.batchNumMessages, Settings.compressionCodec} = do
   let properties =
         Producer.brokersList brokerAddresses
           ++ Producer.sendTimeout deliveryTimeout
           ++ Producer.logLevel logLevel
-          ++ Producer.compression Producer.Snappy
+          ++ Producer.compression compressionCodec
           ++ Producer.extraProps
             ( Dict.fromList
                 [ ( "batch.num.messages",
@@ -263,8 +283,8 @@
             record'
             ( \deliveryReport -> do
                 log <- Platform.silentHandler
-                Task.perform log
-                  <| case deliveryReport of
+                Task.perform log <|
+                  case deliveryReport of
                     Producer.DeliverySuccess _producerRecord _offset -> onDeliveryCallback
                     _ -> Task.succeed ()
             )
diff --git a/src/Kafka/Settings.hs b/src/Kafka/Settings.hs
--- a/src/Kafka/Settings.hs
+++ b/src/Kafka/Settings.hs
@@ -21,7 +21,9 @@
     -- | Message delivery timeout. See hw-kafka's documentation for more info
     deliveryTimeout :: Kafka.Producer.Timeout,
     -- | Number of messages to batch together before sending to Kafka.
-    batchNumMessages :: BatchNumMessages
+    batchNumMessages :: BatchNumMessages,
+    -- | Compression codec used for topics
+    compressionCodec :: Internal.KafkaCompressionCodec
   }
 
 -- | Number of messages to batch together before sending to Kafka.
@@ -38,12 +40,13 @@
 -- KAFKA_BATCH_SIZE=10000
 decoder :: Environment.Decoder Settings
 decoder =
-  map4
+  map5
     Settings
     Internal.decoderBrokerAddresses
     Internal.decoderKafkaLogLevel
     decoderDeliveryTimeout
     decoderBatchNumMessages
+    Internal.decoderCompressionCodec
 
 decoderDeliveryTimeout :: Environment.Decoder Kafka.Producer.Timeout
 decoderDeliveryTimeout =
diff --git a/src/Kafka/Settings/Internal.hs b/src/Kafka/Settings/Internal.hs
--- a/src/Kafka/Settings/Internal.hs
+++ b/src/Kafka/Settings/Internal.hs
@@ -1,7 +1,9 @@
 module Kafka.Settings.Internal
   ( Kafka.Types.KafkaLogLevel (..),
+    Kafka.Types.KafkaCompressionCodec (..),
     decoderBrokerAddresses,
     decoderKafkaLogLevel,
+    decoderCompressionCodec,
   )
 where
 
@@ -26,17 +28,30 @@
         Environment.description = "Kafka log level",
         Environment.defaultValue = "Debug"
       }
-    (map kafkaLogLevelFromText Environment.text)
+    ( Environment.enum
+        [ ("Emerg", Kafka.Types.KafkaLogEmerg),
+          ("Alert", Kafka.Types.KafkaLogAlert),
+          ("Crit", Kafka.Types.KafkaLogCrit),
+          ("Err", Kafka.Types.KafkaLogErr),
+          ("Warning", Kafka.Types.KafkaLogWarning),
+          ("Notice", Kafka.Types.KafkaLogNotice),
+          ("Info", Kafka.Types.KafkaLogInfo),
+          ("Debug", Kafka.Types.KafkaLogDebug)
+        ]
+    )
 
-kafkaLogLevelFromText :: Text -> Kafka.Types.KafkaLogLevel
-kafkaLogLevelFromText text =
-  case text of
-    "Emerg" -> Kafka.Types.KafkaLogEmerg
-    "Alert" -> Kafka.Types.KafkaLogAlert
-    "Crit" -> Kafka.Types.KafkaLogCrit
-    "Err" -> Kafka.Types.KafkaLogErr
-    "Warning" -> Kafka.Types.KafkaLogWarning
-    "Notice" -> Kafka.Types.KafkaLogNotice
-    "Info" -> Kafka.Types.KafkaLogInfo
-    "Debug" -> Kafka.Types.KafkaLogDebug
-    _ -> Kafka.Types.KafkaLogDebug
+decoderCompressionCodec :: Environment.Decoder Kafka.Types.KafkaCompressionCodec
+decoderCompressionCodec =
+  Environment.variable
+    Environment.Variable
+      { Environment.name = "KAFKA_COMPRESSION_CODEC",
+        Environment.description = "Compression codec used for topics. Supported values are: NoCopmression, Gzip, Snappy and Lz4",
+        Environment.defaultValue = "Snappy"
+      }
+    ( Environment.enum
+        [ ("NoCompression", Kafka.Types.NoCompression),
+          ("Gzip", Kafka.Types.Gzip),
+          ("Snappy", Kafka.Types.Snappy),
+          ("Lz4", Kafka.Types.Lz4)
+        ]
+    )
diff --git a/src/Kafka/Test.hs b/src/Kafka/Test.hs
--- a/src/Kafka/Test.hs
+++ b/src/Kafka/Test.hs
@@ -15,7 +15,7 @@
 -- yields a mock Kafka handler, and returns an expectation wrapping a list of
 -- messages that would have been written if the handler was real
 stub ::
-  Stack.HasCallStack =>
+  (Stack.HasCallStack) =>
   (Internal.Handler -> Expect.Expectation) ->
   Expect.Expectation' (List Kafka.Msg)
 stub stubbed = do
diff --git a/src/Kafka/Worker.hs b/src/Kafka/Worker.hs
--- a/src/Kafka/Worker.hs
+++ b/src/Kafka/Worker.hs
@@ -24,8 +24,10 @@
     Internal.TopicSubscription,
     Internal.subscription,
     Internal.subscriptionManageOwnOffsets,
+    Internal.receiveRawMessages,
     Internal.PartitionOffset (..),
     Partition.SeekCmd (..),
+    Partition.ProcessAttemptsCount (..),
     Internal.CommitToKafkaAsWell (..),
   )
 where
diff --git a/src/Kafka/Worker/Fetcher.hs b/src/Kafka/Worker/Fetcher.hs
--- a/src/Kafka/Worker/Fetcher.hs
+++ b/src/Kafka/Worker/Fetcher.hs
@@ -1,6 +1,7 @@
 module Kafka.Worker.Fetcher (pollingLoop) where
 
 import qualified Control.Concurrent
+import qualified Control.Concurrent.MVar as MVar
 import qualified Control.Exception.Safe as Exception
 import qualified Data.ByteString as ByteString
 import qualified Dict
@@ -23,10 +24,11 @@
   EnqueueRecord ->
   Analytics.Analytics ->
   Consumer.KafkaConsumer ->
+  MVar.MVar () ->
   Prelude.IO ()
-pollingLoop settings enqueueRecord analytics consumer = do
+pollingLoop settings enqueueRecord analytics consumer consumerLock = do
   now <- nextPollingTimestamp
-  pollingLoop' settings enqueueRecord analytics consumer (pollTimeIsOld now)
+  pollingLoop' settings enqueueRecord analytics consumer consumerLock (pollTimeIsOld now)
 
 newtype LastPollingTimestamp = LastPollingTimestamp Float
 
@@ -46,6 +48,7 @@
   EnqueueRecord ->
   Analytics.Analytics ->
   Consumer.KafkaConsumer ->
+  MVar.MVar () ->
   LastPollingTimestamp ->
   Prelude.IO ()
 pollingLoop'
@@ -58,10 +61,23 @@
   enqueueRecord
   analytics
   consumer
+  consumerLock
   lastPollTimestamp = do
     -- we block here if we're actively revoking
     -- Check whether we need to shut down while long-polling for new messages.
-    eitherMsgs <- Consumer.pollMessageBatch consumer pollingTimeout pollBatchSize
+    eitherMsgs <-
+      -- We use a lock to prevent running this concurrently with pause/resume calls, due to bugs in
+      -- librdkafka, fixed in 2.1.0, while hw-kafka is on 1.6. Search Worker/Internal.hs for
+      -- consumerLock for the other side of this.
+      --
+      -- The symptom is messages being skipped every once in a while in a slow consumer that has its
+      -- buffer filled up and had to pause/resume all the time.
+      --
+      -- See https://github.com/confluentinc/librdkafka/blob/c282ba2423b2694052393c8edb0399a5ef471b3f/CHANGELOG.md?plain=1#L90-L95
+      --
+      -- We have a small app to reproduce the bug. Check out scripts/pause-resume-bug/README.md
+      MVar.withMVar consumerLock <|
+        \_ -> Consumer.pollMessageBatch consumer pollingTimeout pollBatchSize
     msgs <- Prelude.traverse handleKafkaError eitherMsgs
     assignment <-
       Consumer.assignment consumer
@@ -85,7 +101,7 @@
       |> seek consumer
     now <- nextPollingTimestamp
     throttle maxMsgsPerSecondPerPartition maxPollIntervalMs (List.length appendResults) analytics now lastPollTimestamp
-    pollingLoop' settings enqueueRecord analytics consumer (pollTimeIsOld now)
+    pollingLoop' settings enqueueRecord analytics consumer consumerLock (pollTimeIsOld now)
 
 getPartitionKey :: Consumer.ConsumerRecord k v -> (Consumer.TopicName, Consumer.PartitionId)
 getPartitionKey record =
@@ -192,7 +208,7 @@
     then Control.Concurrent.threadDelay microSecondsToSleep
     else Prelude.pure ()
 
-groupDictAndMap :: Ord b => (a -> (b, c)) -> List a -> Dict.Dict b (List c)
+groupDictAndMap :: (Ord b) => (a -> (b, c)) -> List a -> Dict.Dict b (List c)
 groupDictAndMap f =
   List.foldr
     ( \x ->
diff --git a/src/Kafka/Worker/Internal.hs b/src/Kafka/Worker/Internal.hs
--- a/src/Kafka/Worker/Internal.hs
+++ b/src/Kafka/Worker/Internal.hs
@@ -5,9 +5,11 @@
 import qualified Conduit
 import qualified Control.Concurrent
 import qualified Control.Concurrent.Async as Async
+import qualified Control.Concurrent.MVar as MVar
 import qualified Control.Concurrent.STM as STM
 import qualified Control.Concurrent.STM.TVar as TVar
 import qualified Control.Exception.Safe as Exception
+import Control.Monad (unless)
 import qualified Data.Aeson as Aeson
 import qualified Data.UUID
 import qualified Data.UUID.V4
@@ -53,7 +55,8 @@
   { topic :: Kafka.Topic,
     onMessage :: Partition.MessageCallback,
     offsetSource :: OffsetSource,
-    commitToKafkaAsWell :: CommitToKafkaAsWell
+    commitToKafkaAsWell :: CommitToKafkaAsWell,
+    messageFormat :: Partition.MessageFormat
   }
 
 -- | Commit the offset to Kafka in addition to an externally managed storage.
@@ -68,6 +71,7 @@
     -- | The partition's offset.
     offset :: Int
   }
+  deriving (Show)
 
 -- | Create a subscription for a topic.
 --
@@ -89,12 +93,13 @@
     { topic = Kafka.Topic topic,
       onMessage =
         Partition.MessageCallback
-          ( \_ msg -> do
+          ( \_ _ msg -> do
               callback msg
               Task.succeed Partition.NoSeek
           ),
       offsetSource = InKafka,
-      commitToKafkaAsWell = CommitToKafkaAsWell
+      commitToKafkaAsWell = CommitToKafkaAsWell,
+      messageFormat = Partition.WithMetadata
     }
 
 -- | Create a subscription for a topic and manage offsets for that topic
@@ -114,14 +119,14 @@
 -- >              sql
 -- >                "SELECT partition, offset FROM offsets WHERE partition = %"
 -- >                [partitions] )
--- >           (\msg -> Debug.todo "Process your message here!")
+-- >           (\retryCount msg -> Debug.todo "Process your message here!")
 -- >   process settings subscription
 subscriptionManageOwnOffsets ::
   (Aeson.FromJSON msg, Aeson.ToJSON msg) =>
   Text ->
   CommitToKafkaAsWell ->
   ([Int] -> Task Text (List PartitionOffset)) ->
-  (PartitionOffset -> msg -> Task Text Partition.SeekCmd) ->
+  (PartitionOffset -> Partition.ProcessAttemptsCount -> msg -> Task Text Partition.SeekCmd) ->
   TopicSubscription
 subscriptionManageOwnOffsets topic commitToKafkaAsWell fetchOffsets callback =
   TopicSubscription
@@ -129,7 +134,7 @@
       commitToKafkaAsWell,
       onMessage =
         Partition.MessageCallback
-          ( \record msg -> do
+          ( \record retryCount msg -> do
               let offsetParams =
                     PartitionOffset
                       { partitionId =
@@ -137,7 +142,7 @@
                             |> partitionIdToInt,
                         offset = Consumer.unOffset (Consumer.crOffset record)
                       }
-              callback offsetParams msg
+              callback offsetParams retryCount msg
           ),
       offsetSource =
         Elsewhere
@@ -149,7 +154,8 @@
               offsets
                 |> List.map toPartitionKey
                 |> Task.succeed
-          )
+          ),
+      messageFormat = Partition.WithMetadata
     }
   where
     toPartitionKey :: PartitionOffset -> (PartitionKey, Int)
@@ -163,6 +169,12 @@
     partitionIdToInt :: Consumer.PartitionId -> Int
     partitionIdToInt (Consumer.PartitionId int) = Prelude.fromIntegral int
 
+-- | Modify the subscription to accept raw messages without metadata.
+-- This is useful to receive messages generated by other services.
+receiveRawMessages :: TopicSubscription -> TopicSubscription
+receiveRawMessages topicSubscription =
+  topicSubscription {messageFormat = Partition.Raw}
+
 -- | This determines how a worker that was just assigned a partition should
 -- decide at which message offset to continue processing.
 data OffsetSource where
@@ -179,33 +191,36 @@
 -- | Starts the kafka worker handling messages.
 process :: Settings.Settings -> Text -> TopicSubscription -> Prelude.IO ()
 process settings groupIdText topicSubscriptions = do
-  processWithoutShutdownEnsurance settings (Consumer.ConsumerGroupId groupIdText) topicSubscriptions
-  -- Start an ensurance policy to make sure we exit in 5 seconds. We've seen
-  -- cases where our graceful shutdown seems to hang, resulting in a worker
-  -- that's not doing anything. We should try to fix those failures, but for the
-  -- ones that remain this is our fallback.
-  --
-  -- Running it using `Async.async` makes it so we won't wait for this thread to
-  -- complete. If the regular shutdown completes before this thread is done we
-  -- will exit early.
-  _ <-
-    Async.async <| do
-      Control.Concurrent.threadDelay 5_000_000 {- 5 seconds -}
-      Prelude.putStrLn "Something is holding up shutdown. Going to die ungracefully now."
-      System.Posix.Process.exitImmediately (System.Exit.ExitFailure 1)
-  Prelude.pure ()
+  state <- initState
+  Async.race_
+    ( do
+        Stopping.waitUntilStopping (stopping state)
 
+        -- Start an ensurance policy to make sure we exit in 5 seconds. We've seen
+        -- cases where our graceful shutdown seems to hang, resulting in a worker
+        -- that's not doing anything. We should try to fix those failures, but for the
+        -- ones that remain this is our fallback.
+        Control.Concurrent.threadDelay 5_000_000 {- 5 seconds -}
+        putTextLn "Something is holding up shutdown. Going to die ungracefully now."
+        System.Posix.Process.exitImmediately (System.Exit.ExitFailure 1)
+    )
+    (process' state settings (Consumer.ConsumerGroupId groupIdText) topicSubscriptions)
+
 -- | Like `process`, but doesn't exit the current process by itself. This risks
 -- leaving zombie processes when used in production but is safer in tests, where
 -- the worker shares the OS process with other test code and the test runner.
 processWithoutShutdownEnsurance :: Settings.Settings -> Consumer.ConsumerGroupId -> TopicSubscription -> Prelude.IO ()
 processWithoutShutdownEnsurance settings groupId topicSubscriptions = do
-  let TopicSubscription {onMessage, topic, offsetSource, commitToKafkaAsWell} = topicSubscriptions
   state <- initState
+  process' state settings groupId topicSubscriptions
+
+process' :: State -> Settings.Settings -> Consumer.ConsumerGroupId -> TopicSubscription -> Prelude.IO ()
+process' state settings groupId topicSubscriptions = do
+  let TopicSubscription {onMessage, topic, offsetSource, commitToKafkaAsWell, messageFormat} = topicSubscriptions
   onQuitSignal (Stopping.stopTakingRequests (stopping state) "Received stop signal")
   Conduit.withAcquire (Observability.handler (Settings.observability settings)) <| \observabilityHandler -> do
     Exception.bracketWithError
-      (createConsumer settings groupId observabilityHandler offsetSource commitToKafkaAsWell onMessage topic state)
+      (createConsumer settings groupId observabilityHandler offsetSource commitToKafkaAsWell onMessage topic messageFormat state)
       (cleanUp observabilityHandler (rebalanceInfo state) (stopping state))
       (runThreads settings state)
 
@@ -233,6 +248,7 @@
   CommitToKafkaAsWell ->
   Partition.MessageCallback ->
   Kafka.Topic ->
+  Partition.MessageFormat ->
   State ->
   Prelude.IO Consumer.KafkaConsumer
 createConsumer
@@ -240,7 +256,8 @@
     { Settings.brokerAddresses,
       Settings.logLevel,
       Settings.maxPollIntervalMs,
-      Settings.onProcessMessageSkip
+      Settings.onProcessMessageSkip,
+      Settings.compressionCodec
     }
   groupId
   observability
@@ -248,10 +265,12 @@
   commitToKafkaAsWell
   callback
   topic
+  messageFormat
   state = do
     let rebalance =
           rebalanceCallback
             onProcessMessageSkip
+            messageFormat
             observability
             callback
             offsetSource
@@ -264,7 +283,7 @@
             ++ Consumer.noAutoCommit
             ++ Consumer.logLevel logLevel
             ++ Consumer.setCallback (Consumer.rebalanceCallback rebalance)
-            ++ Consumer.compression Consumer.Snappy
+            ++ Consumer.compression compressionCodec
             ++ Consumer.extraProps
               ( Dict.fromList
                   [("max.poll.interval.ms", Text.fromInt (Settings.unMaxPollIntervalMs maxPollIntervalMs))]
@@ -293,6 +312,7 @@
 -- https://docs.confluent.io/2.0.0/clients/librdkafka/classRdKafka_1_1RebalanceCb.html#a490a91c52724382a72380af621958741
 rebalanceCallback ::
   Settings.SkipOrNot ->
+  Partition.MessageFormat ->
   Observability.Handler ->
   Partition.MessageCallback ->
   OffsetSource ->
@@ -301,7 +321,7 @@
   Consumer.KafkaConsumer ->
   Consumer.RebalanceEvent ->
   Prelude.IO ()
-rebalanceCallback skipOrNot observability callback offsetSource commitToKafkaAsWell state consumer rebalanceEvent = do
+rebalanceCallback skipOrNot messageFormat observability callback offsetSource commitToKafkaAsWell state consumer rebalanceEvent = do
   now <- GHC.Clock.getMonotonicTime
   Analytics.updateTimeOfLastRebalance now (analytics state)
   case rebalanceEvent of
@@ -366,14 +386,15 @@
           ( \(partitionKey, offset) -> do
               initPartition
                 skipOrNot
+                messageFormat
                 offset
                 observability
                 consumer
                 callback
                 state
                 partitionKey
-              STM.atomically
-                <| TVar.modifyTVar' (rebalanceInfo state) (Dict.insert partitionKey (Assign, now))
+              STM.atomically <|
+                TVar.modifyTVar' (rebalanceInfo state) (Dict.insert partitionKey (Assign, now))
           )
         |> map (\_ -> ())
     Consumer.RebalanceAssign _ -> Prelude.pure ()
@@ -425,7 +446,7 @@
 -- | Disconnects our Consumer / yields back partitions on quit / node shutdown
 cleanUp :: Observability.Handler -> RebalanceInfo -> Stopping.Stopping -> Maybe Exception.SomeException -> Consumer.KafkaConsumer -> Prelude.IO ()
 cleanUp observabilityHandler rebalanceInfo stopping maybeException consumer = do
-  Prelude.putStrLn "Cleaning up"
+  putTextLn "Cleaning up"
   _ <- Consumer.closeConsumer consumer
   -- In case we're already stopping, get the reason we're doing so.
   maybeStopReason <- Stopping.stopReason stopping
@@ -435,6 +456,7 @@
   -- at some point, k8s should report system crashes. In the mean time, we'll do it.
   Platform.rootTracingSpanIO
     requestId
+    Platform.silentTrack
     (Observability.report observabilityHandler requestId)
     "Kafka consumer shutting down"
     <| \log -> do
@@ -451,9 +473,9 @@
             |> Task.perform log
   writeCrashLogOnError maybeException
   case (maybeException, maybeStopReason) of
-    (Just exception, _) -> Prelude.putStrLn ("Shut down because of exception: " ++ Exception.displayException exception)
-    (_, Just stopReason) -> Prelude.putStrLn ("Shut down because of: " ++ Text.toList stopReason)
-    (Nothing, Nothing) -> Prelude.putStrLn "Shut down for an unknown reason."
+    (Just exception, _) -> putTextLn ("Shut down because of exception: " ++ Text.fromList (Exception.displayException exception))
+    (_, Just stopReason) -> putTextLn ("Shut down because of: " ++ stopReason)
+    (Nothing, Nothing) -> putTextLn "Shut down for an unknown reason."
 
 -- | Handle crash logging
 writeCrashLogOnError :: Maybe Exception.SomeException -> Prelude.IO ()
@@ -477,6 +499,7 @@
 -- (after the Consumer.rebalanceassign event occurs)
 initPartition ::
   Settings.SkipOrNot ->
+  Partition.MessageFormat ->
   Partition.CommitOffsets ->
   Observability.Handler ->
   Consumer.KafkaConsumer ->
@@ -484,10 +507,11 @@
   State ->
   PartitionKey ->
   Prelude.IO ()
-initPartition skipOrNot commitOffset observabilityHandler consumer callback state key = do
+initPartition skipOrNot messageFormat commitOffset observabilityHandler consumer callback state key = do
   -- # Start worker thread for handling messages in partition.
   Partition.spawnWorkerThread
     skipOrNot
+    messageFormat
     commitOffset
     observabilityHandler
     (analytics state)
@@ -517,7 +541,7 @@
         ( do
             -- Remove the partition from the dict to clean up memory
             STM.atomically <| TVar.modifyTVar' (partitions state) (Dict.remove key)
-            Prelude.putStrLn ("Stop processing messages for partition: " ++ Prelude.show key)
+            putTextLn ("Stop processing messages for partition: " ++ Text.tshow key)
         )
     )
 
@@ -527,12 +551,13 @@
   Consumer.KafkaConsumer ->
   Prelude.IO ()
 runThreads settings state consumer = do
+  consumerLock <- MVar.newMVar ()
   Stopping.runUnlessStopping
     (stopping state)
     ()
     ( Async.race
-        (pauseAndAnalyticsLoop (Settings.maxMsgsPerPartitionBufferedLocally settings) consumer state Set.empty)
-        (Fetcher.pollingLoop settings (enqueueRecord (partitions state)) (analytics state) consumer)
+        (pauseAndAnalyticsLoop (Settings.maxMsgsPerPartitionBufferedLocally settings) consumer consumerLock state Set.empty)
+        (Fetcher.pollingLoop settings (enqueueRecord (partitions state)) (analytics state) consumer consumerLock)
         |> map (\_ -> ())
     )
 
@@ -561,18 +586,33 @@
 pauseAndAnalyticsLoop ::
   Settings.MaxMsgsPerPartitionBufferedLocally ->
   Consumer.KafkaConsumer ->
+  MVar.MVar () ->
   State ->
   Set.Set PartitionKey ->
   Prelude.IO ()
-pauseAndAnalyticsLoop maxBufferSize consumer state pausedPartitions = do
+pauseAndAnalyticsLoop maxBufferSize consumer consumerLock state pausedPartitions = do
   desiredPausedPartitions <- pausedPartitionKeys maxBufferSize (partitions state)
   Analytics.updatePaused (Set.size desiredPausedPartitions) (analytics state)
   let newlyPaused = Set.diff desiredPausedPartitions pausedPartitions
-  _ <- Consumer.pausePartitions consumer (Set.toList newlyPaused)
   let newlyResumed = Set.diff pausedPartitions desiredPausedPartitions
-  _ <- Consumer.resumePartitions consumer (Set.toList newlyResumed)
+  -- We use a lock to prevent running this concurrently with pollMessageBatch calls, due to bugs in
+  -- librdkafka, fixed in 2.1.0, while hw-kafka is on 1.6. Search Worker/Fetcher.hs for consumerLock
+  -- for the other side of this.
+  --
+  -- The symptom is messages being skipped every once in a while in a slow consumer that has its
+  -- buffer filled up and has to pause/resume all the time.
+  --
+  -- See https://github.com/confluentinc/librdkafka/blob/c282ba2423b2694052393c8edb0399a5ef471b3f/CHANGELOG.md?plain=1#L90-L95
+  --
+  -- We have a small app to reproduce the bug. Check out scripts/pause-resume-bug/README.md
+  unless (Set.isEmpty newlyPaused && Set.isEmpty newlyResumed) <|
+    MVar.withMVar consumerLock <|
+      \_ -> do
+        _ <- Consumer.pausePartitions consumer (Set.toList newlyPaused)
+        _ <- Consumer.resumePartitions consumer (Set.toList newlyResumed)
+        Prelude.pure ()
   Control.Concurrent.threadDelay 1_000_000 {- 1 second -}
-  pauseAndAnalyticsLoop maxBufferSize consumer state desiredPausedPartitions
+  pauseAndAnalyticsLoop maxBufferSize consumer consumerLock state desiredPausedPartitions
 
 pausedPartitionKeys :: Settings.MaxMsgsPerPartitionBufferedLocally -> AllPartitions -> Prelude.IO (Set.Set PartitionKey)
 pausedPartitionKeys (Settings.MaxMsgsPerPartitionBufferedLocally maxBufferSize) partitions = do
@@ -582,8 +622,8 @@
     |> Prelude.traverse
       ( \(key, partition) -> do
           maybeLen <- Partition.length partition
-          Prelude.pure
-            <| case maybeLen of
+          Prelude.pure <|
+            case maybeLen of
               Nothing -> Nothing
               Just length ->
                 if length > maxBufferSize
diff --git a/src/Kafka/Worker/Partition.hs b/src/Kafka/Worker/Partition.hs
--- a/src/Kafka/Worker/Partition.hs
+++ b/src/Kafka/Worker/Partition.hs
@@ -10,6 +10,8 @@
     MessageCallback (..),
     SeekCmd (..),
     CommitOffsets (..),
+    MessageFormat (..),
+    ProcessAttemptsCount (..),
     -- just exported for tests
     microSecondsDelayForAttempt,
     OnStartup (OnStartup),
@@ -94,7 +96,7 @@
 data MessageCallback where
   MessageCallback ::
     (Show e, Aeson.ToJSON msg, Aeson.FromJSON msg) =>
-    (Consumer.ConsumerRecord () () -> msg -> Task e SeekCmd) ->
+    (Consumer.ConsumerRecord () () -> ProcessAttemptsCount -> msg -> Task e SeekCmd) ->
     MessageCallback
 
 data CommitOffsets
@@ -107,6 +109,11 @@
     -- out of date.
     ElsewhereButToKafkaAsWell Int
 
+-- | Expected format of the messages.
+data MessageFormat
+  = WithMetadata
+  | Raw
+
 -- | A thread that processes messages for a particular partition. Cleans itself
 -- up if it ever runs out.
 --
@@ -139,6 +146,7 @@
 -- such a scheme. As we run this code we'll gather data that can help us decide.
 spawnWorkerThread ::
   Settings.SkipOrNot ->
+  MessageFormat ->
   CommitOffsets ->
   Observability.Handler ->
   Analytics.Analytics ->
@@ -148,20 +156,21 @@
   OnStartup ->
   OnCleanup ->
   Prelude.IO ()
-spawnWorkerThread skipOrNot commitOffsets observabilityHandler analytics stopping consumer callback (OnStartup onStartup) (OnCleanup onCleanup) = do
+spawnWorkerThread skipOrNot messageFormat commitOffsets observabilityHandler analytics stopping consumer callback (OnStartup onStartup) (OnCleanup onCleanup) = do
   -- Synchronously create the queue that will come to contain messages for the
   -- partition. This way we'll be able to start receiving messages for this
   -- partition as soon as this function returns, even if the processing thread
   -- we start below still needs boot.
   partition <-
-    map Partition <| TVar.newTVarIO
-      <| case commitOffsets of
-        ToKafka -> Assigned Seq.empty
-        Elsewhere offset -> AwaitingSeekTo offset
-        ElsewhereButToKafkaAsWell offset -> AwaitingSeekTo offset
+    map Partition <|
+      TVar.newTVarIO <|
+        case commitOffsets of
+          ToKafka -> Assigned Seq.empty
+          Elsewhere offset -> AwaitingSeekTo offset
+          ElsewhereButToKafkaAsWell offset -> AwaitingSeekTo offset
   onStartup partition
   Exception.finally
-    (processMsgLoop skipOrNot commitOffsets observabilityHandler State {analytics, stopping, partition} consumer callback)
+    (processMsgLoop skipOrNot messageFormat commitOffsets observabilityHandler State {analytics, stopping, partition} consumer callback)
     onCleanup
     |> Async.async
     -- If the async process spawned here throws an exception, rethrow it
@@ -182,29 +191,20 @@
 
 processMsgLoop ::
   Settings.SkipOrNot ->
+  MessageFormat ->
   CommitOffsets ->
   Observability.Handler ->
   State ->
   Consumer.KafkaConsumer ->
   MessageCallback ->
   Prelude.IO ()
-processMsgLoop skipOrNot commitOffsets observabilityHandler state consumer callback@(MessageCallback runCallback) = do
+processMsgLoop skipOrNot messageFormat commitOffsets observabilityHandler state consumer callback@(MessageCallback runCallback) = do
   -- # Get the next message from the queue.
   peekResponse <- peekRecord state
   case peekResponse of
     StopThread ->
       Prelude.pure ()
     (NextMsg processAttempts record) -> do
-      case processAttempts of
-        (ProcessAttemptsCount 0) -> Prelude.pure ()
-        (ProcessAttemptsCount attempts) ->
-          -- Wait a bit if this is a retry, to prevent putting a lot of retry
-          -- stress on downstream systems or generating huge numbers of error
-          -- messages.
-          microSecondsDelayForAttempt attempts
-            |> Prelude.fromIntegral
-            |> Control.Concurrent.threadDelay
-
       doAnything <- Platform.doAnythingHandler
       let commit processResult =
             case processResult of
@@ -248,6 +248,7 @@
       (RequestId requestId, details) <- getTracingDetails (analytics state) processAttempts record
       Platform.rootTracingSpanIO
         requestId
+        Platform.silentTrack
         (Observability.report observabilityHandler requestId)
         "Assigned Kafka message"
         ( \log -> do
@@ -255,8 +256,8 @@
             -- at least we'll have nice context in logs!
             Platform.setTracingSpanDetailsIO log details
             handleFailures log <| do
-              msg <- decodeMessage record
-              runCallback record {Consumer.crKey = (), Consumer.crValue = ()} msg
+              msg <- decodeMessage messageFormat record
+              runCallback record {Consumer.crKey = (), Consumer.crValue = ()} processAttempts msg
                 |> Task.mapError WorkerCallbackFailed
                 |> Task.onError
                   ( \err -> do
@@ -270,6 +271,7 @@
       -- # Loop for the next message
       processMsgLoop
         skipOrNot
+        messageFormat
         commitOffsets
         observabilityHandler
         state
@@ -278,12 +280,11 @@
 
 microSecondsDelayForAttempt :: Int -> Int
 microSecondsDelayForAttempt attempts =
-  min
-    3_600_000_000 {- 1 hour in microseconds -}
-    ((10 Prelude.^ attempts) * 1000_000 {- 1 second in microseconds -})
+  -- Maximum delay is 2^10 seconds = 17 minutes
+  2 Prelude.^ (min attempts 10) * 1_000_000
 
 handleFailures ::
-  Show e =>
+  (Show e) =>
   Platform.LogHandler ->
   Task (WorkerError e) a ->
   Prelude.IO ()
@@ -360,6 +361,7 @@
                       Prelude.Left _ -> Nothing
                       Prelude.Right keyText -> Just keyText
                 ),
+          Log.Kafka.offset = Just (Consumer.unOffset (Consumer.crOffset record)),
           Log.Kafka.contents = Just contents,
           Log.Kafka.processAttempt = Just processAttempt,
           Log.Kafka.createTime,
@@ -374,26 +376,38 @@
 millisToSecs :: Consumer.Millis -> Clock.UTCTime
 millisToSecs (Consumer.Millis millis) = fromPosix (millis // 1000)
 
-decodeMessage :: (Aeson.FromJSON msg) => ConsumerRecord -> Task (WorkerError e) msg
-decodeMessage record = do
-  let eitherMsg =
+decodeMessage :: (Aeson.FromJSON msg) => MessageFormat -> ConsumerRecord -> Task (WorkerError e) msg
+decodeMessage messageFormat record = do
+  let json =
         Consumer.crValue record
           -- We'll accept the absence of a message if the worker expects a message
           -- of type `()`. The default JSON encoding for `()` is "[]".
           |> Maybe.withDefault "[]"
-          |> Aeson.eitherDecodeStrict
-  case eitherMsg of
-    Prelude.Left err ->
-      Task.fail (MsgDecodingFailed (Text.fromList err))
-    Prelude.Right msgWithMetaData ->
-      case Internal.value msgWithMetaData of
-        (Internal.Encodable value) ->
-          case Aeson.fromJSON (Aeson.toJSON value) of
-            Aeson.Error err ->
-              Task.fail (MsgDecodingFailed (Text.fromList err))
-            Aeson.Success msg ->
-              Task.succeed msg
 
+      parseMessage ::
+        forall msg a e.
+        (Aeson.FromJSON a) =>
+        (a -> Task (WorkerError e) msg) ->
+        Task (WorkerError e) msg
+      parseMessage = \fn ->
+        case json |> Aeson.eitherDecodeStrict of
+          Prelude.Left err ->
+            Task.fail (MsgDecodingFailed (Text.fromList err))
+          Prelude.Right m -> fn m
+
+  case messageFormat of
+    WithMetadata ->
+      parseMessage <| \msgWithMetaData ->
+        case Internal.value msgWithMetaData of
+          (Internal.Encodable value) ->
+            case Aeson.fromJSON (Aeson.toJSON value) of
+              Aeson.Error err ->
+                Task.fail (MsgDecodingFailed (Text.fromList err))
+              Aeson.Success msg ->
+                Task.succeed msg
+    Raw ->
+      parseMessage Task.succeed
+
 commitRecord ::
   Platform.DoAnythingHandler ->
   Consumer.KafkaConsumer ->
@@ -426,24 +440,41 @@
   Stopping.runUnlessStopping
     (stopping state)
     StopThread
-    ( STM.atomically
-        <| do
-          let (Partition partition') = partition state
-          backlog' <- TVar.readTVar partition'
-          case backlog' of
-            AwaitingSeekTo _ ->
-              STM.retry
-            Stopping -> do
-              Prelude.pure StopThread
-            Assigned Seq.Empty ->
-              STM.retry
-            Assigned ((processAttemptsCount, first) Seq.:<| rest) -> do
-              -- Bump the retry count so that the next time we read this message, we
-              -- know we've read it before.
-              TVar.writeTVar
-                partition'
-                (Assigned ((processAttemptsCount + 1, first) Seq.:<| rest))
-              Prelude.pure (NextMsg processAttemptsCount first)
+    ( do
+        next <-
+          STM.atomically <|
+            do
+              let (Partition partition') = partition state
+              backlog' <- TVar.readTVar partition'
+              case backlog' of
+                AwaitingSeekTo _ ->
+                  STM.retry
+                Stopping -> do
+                  Prelude.pure Nothing
+                Assigned Seq.Empty ->
+                  STM.retry
+                Assigned ((processAttemptsCount, first) Seq.:<| rest) -> do
+                  -- Bump the retry count so that the next time we read this message, we
+                  -- know we've read it before.
+                  TVar.writeTVar
+                    partition'
+                    (Assigned ((processAttemptsCount + 1, first) Seq.:<| rest))
+                  Prelude.pure <| Just (processAttemptsCount, first)
+
+        case next of
+          Nothing -> Prelude.pure StopThread
+          Just (processAttempts, record) -> do
+            case processAttempts of
+              (ProcessAttemptsCount 0) -> Prelude.pure ()
+              (ProcessAttemptsCount attempts) ->
+                -- Wait a bit if this is a retry, to prevent putting a lot of retry
+                -- stress on downstream systems or generating huge numbers of error
+                -- messages.
+                microSecondsDelayForAttempt attempts
+                  |> Prelude.fromIntegral
+                  |> Control.Concurrent.threadDelay
+
+            Prelude.pure (NextMsg processAttempts record)
     )
 
 awaitingSeekTo :: Partition -> Int -> Prelude.IO ()
@@ -499,8 +530,8 @@
 length :: Partition -> Prelude.IO (Maybe Int)
 length (Partition partition) = do
   backlog <- TVar.readTVarIO partition
-  Prelude.pure
-    <| case backlog of
+  Prelude.pure <|
+    case backlog of
       AwaitingSeekTo _ -> Nothing
       Stopping -> Nothing
       Assigned queue ->
diff --git a/src/Kafka/Worker/Settings.hs b/src/Kafka/Worker/Settings.hs
--- a/src/Kafka/Worker/Settings.hs
+++ b/src/Kafka/Worker/Settings.hs
@@ -33,7 +33,9 @@
     maxPollIntervalMs :: MaxPollIntervalMs,
     -- | This option provides us the possibility to skip messages on failure.
     -- Useful for testing Kafka worker. DoNotSkip is a reasonable default!
-    onProcessMessageSkip :: SkipOrNot
+    onProcessMessageSkip :: SkipOrNot,
+    -- | Compression codec used for topics
+    compressionCodec :: Internal.KafkaCompressionCodec
   }
 
 -- | This option provides us the possibility to skip messages on failure.
@@ -74,6 +76,7 @@
     |> andMap decoderPollBatchSize
     |> andMap decoderMaxPollIntervalMs
     |> andMap decoderOnProcessMessageFailure
+    |> andMap Internal.decoderCompressionCodec
 
 decoderPollingTimeout :: Environment.Decoder Consumer.Timeout
 decoderPollingTimeout =
diff --git a/src/Kafka/Worker/Stopping.hs b/src/Kafka/Worker/Stopping.hs
--- a/src/Kafka/Worker/Stopping.hs
+++ b/src/Kafka/Worker/Stopping.hs
@@ -3,6 +3,7 @@
     stopReason,
     stopTakingRequests,
     runUnlessStopping,
+    waitUntilStopping,
     Stopping,
   )
 where
@@ -22,14 +23,14 @@
 
 stopTakingRequests :: Stopping -> Text -> Prelude.IO ()
 stopTakingRequests (Stopping stopping) reason = do
-  Prelude.putStrLn "Gracefully shutting down..."
+  putTextLn "Gracefully shutting down..."
   MVar.tryPutMVar stopping reason
     |> map (\_ -> ())
 
 runUnlessStopping :: Stopping -> a -> Prelude.IO a -> Prelude.IO a
-runUnlessStopping (Stopping stopping) stoppingVal action =
+runUnlessStopping stopping stoppingVal action =
   Async.race
-    (MVar.readMVar stopping |> map (\_ -> ()))
+    (waitUntilStopping stopping)
     action
     |> map
       ( \either ->
@@ -37,3 +38,7 @@
             Prelude.Left () -> stoppingVal
             Prelude.Right r -> r
       )
+
+waitUntilStopping :: Stopping -> Prelude.IO ()
+waitUntilStopping (Stopping stopping) =
+  MVar.readMVar stopping |> map (\_ -> ())
diff --git a/test/Helpers.hs b/test/Helpers.hs
--- a/test/Helpers.hs
+++ b/test/Helpers.hs
@@ -4,6 +4,7 @@
     stopWorker,
     test,
     sendSync,
+    spawnWorkerManagingOwnOffsets,
   )
 where
 
@@ -78,6 +79,42 @@
     Async.link async
     Prelude.pure (Worker async)
 
+spawnWorkerManagingOwnOffsets ::
+  (Aeson.ToJSON msg, Aeson.FromJSON msg) =>
+  TestHandler ->
+  Internal.Topic ->
+  (Worker.PartitionOffset -> Worker.ProcessAttemptsCount -> msg -> STM.STM Worker.SeekCmd) ->
+  Expect.Expectation' Worker
+spawnWorkerManagingOwnOffsets handler' topic callback =
+  Expect.fromIO <| do
+    settings <-
+      case Environment.decodeDefaults Worker.Settings.decoder of
+        Ok settings' -> Prelude.pure settings'
+        Err err -> Prelude.fail (Text.toList err)
+    async <-
+      Kafka.Worker.Internal.processWithoutShutdownEnsurance
+        settings
+        (Consumer.ConsumerGroupId "group")
+        ( Worker.subscriptionManageOwnOffsets
+            (Internal.unTopic topic)
+            Worker.CommitToKafkaAsWell
+            ( \partitions ->
+                partitions
+                  |> List.map (\id -> Worker.PartitionOffset {Worker.partitionId = id, Worker.offset = 0})
+                  |> Task.succeed
+            )
+            ( \partitionOffset retryCount msg -> do
+                callback partitionOffset retryCount msg
+                  |> STM.atomically
+                  |> map Ok
+                  |> Platform.doAnything (doAnything handler')
+            )
+        )
+        |> Async.race_ (returnWhenTerminating handler')
+        |> Async.async
+    Async.link async
+    Prelude.pure (Worker async)
+
 -- | Stops a single worker
 stopWorker :: Worker -> Expect.Expectation
 stopWorker (Worker async) =
@@ -86,13 +123,13 @@
 
 -- | creates a test handler
 testHandler :: Settings.Settings -> Prelude.IO TestHandler
-testHandler Settings.Settings {Settings.brokerAddresses, Settings.deliveryTimeout, Settings.logLevel, Settings.batchNumMessages} = do
+testHandler Settings.Settings {Settings.brokerAddresses, Settings.deliveryTimeout, Settings.logLevel, Settings.batchNumMessages, Settings.compressionCodec} = do
   doAnything <- Platform.doAnythingHandler
   let properties =
         Producer.brokersList brokerAddresses
           ++ Producer.sendTimeout deliveryTimeout
           ++ Producer.logLevel logLevel
-          ++ Producer.compression Producer.Snappy
+          ++ Producer.compression compressionCodec
           ++ Producer.extraProps
             ( Dict.fromList
                 [ ( "batch.num.messages",
@@ -119,7 +156,7 @@
       Prelude.pure TestHandler {producer, doAnything, terminator}
 
 -- | puts a message synchronously onto a topic-partition
-sendSync :: Aeson.ToJSON a => TestHandler -> Internal.Topic -> Int -> a -> Expect.Expectation
+sendSync :: (Aeson.ToJSON a) => TestHandler -> Internal.Topic -> Int -> a -> Expect.Expectation
 sendSync handler topicName partitionId msg' =
   Platform.tracingSpan
     "Sync send Kafka messages"
@@ -175,7 +212,7 @@
 
 -- | test helper, that yields a new @Kafka.Topic@ and @TestHandler@
 test ::
-  Stack.HasCallStack =>
+  (Stack.HasCallStack) =>
   Text ->
   ((Internal.Topic, TestHandler) -> Expect.Expectation) ->
   Test.Test
@@ -186,8 +223,8 @@
       ( \task' ->
           Platform.bracketWithError
             ( -- create handler
-              Platform.doAnything doAnything
-                <| case Environment.decodeDefaults Settings.decoder of
+              Platform.doAnything doAnything <|
+                case Environment.decodeDefaults Settings.decoder of
                   Ok settings ->
                     map
                       Ok
diff --git a/test/Main.hs b/test/Main.hs
--- a/test/Main.hs
+++ b/test/Main.hs
@@ -2,11 +2,15 @@
 
 import qualified Spec.Kafka.Worker.Integration
 import qualified Spec.Kafka.Worker.Partition
+import qualified System.Environment
 import qualified Test
 import qualified Prelude
 
 main :: Prelude.IO ()
-main = Test.run tests
+main = do
+  -- macos runners seem to be slow and fail on several kafka integration tests
+  System.Environment.setEnv "NRI_TEST_TIMEOUT" "240000"
+  Test.run tests
 
 tests :: Test.Test
 tests =
diff --git a/test/Spec/Kafka/Worker/Integration.hs b/test/Spec/Kafka/Worker/Integration.hs
--- a/test/Spec/Kafka/Worker/Integration.hs
+++ b/test/Spec/Kafka/Worker/Integration.hs
@@ -4,6 +4,7 @@
 import qualified Dict
 import qualified Expect
 import qualified Helpers
+import qualified Kafka.Worker as Worker
 import qualified Set
 import qualified Test
 import qualified Prelude
@@ -29,7 +30,8 @@
             Helpers.sendSync handler topic 1 (1, 2)
             Helpers.sendSync handler topic 2 (2, 3)
             msgs' <- waitFor msgsTVar (\items -> List.length items == 3)
-            msgs' |> groupDictAndMap identity
+            msgs'
+              |> groupDictAndMap identity
               |> Expect.equal
                 ( Dict.fromList
                     [ (1, [2, 1]),
@@ -54,7 +56,25 @@
                     [ (1, [2, 1]),
                       (2, [2, 1])
                     ]
+                ),
+          Helpers.test "Self-managing workers get retry count info" <| \(topic, handler) -> do
+            Helpers.sendSync handler topic 1 1
+            msgsTVar <- atomically (STM.newTVar Set.empty)
+            _ <-
+              Helpers.spawnWorkerManagingOwnOffsets
+                handler
+                topic
+                ( \partitionOffset (Worker.ProcessAttemptsCount retryCount) msg -> do
+                    STM.modifyTVar' msgsTVar (Set.insert (msg, retryCount))
+                    if retryCount < 1
+                      then STM.throwSTM (Prelude.userError "retry please")
+                      else
+                        Prelude.pure <|
+                          Worker.SeekToOffset ((Worker.offset partitionOffset) + 1)
                 )
+            msgs' <- waitFor msgsTVar (\items -> Set.size items == 1)
+            -- Assert that the message was recorded on its first retry
+            msgs' |> Expect.equal (Set.fromList [(1, 1)])
         ]
     ]
 
@@ -69,7 +89,7 @@
       then Prelude.pure val
       else STM.retry
 
-groupDictAndMap :: Ord b => (a -> (b, c)) -> List a -> Dict.Dict b (List c)
+groupDictAndMap :: (Ord b) => (a -> (b, c)) -> List a -> Dict.Dict b (List c)
 groupDictAndMap f =
   List.foldr
     ( \x ->
diff --git a/test/Spec/Kafka/Worker/Partition.hs b/test/Spec/Kafka/Worker/Partition.hs
--- a/test/Spec/Kafka/Worker/Partition.hs
+++ b/test/Spec/Kafka/Worker/Partition.hs
@@ -12,15 +12,18 @@
         "microSecondsDelayForAttempt"
         [ Test.test "1 attempt" <| \() ->
             Partition.microSecondsDelayForAttempt 1
-              |> Expect.equal 10_000_000,
+              |> Expect.equal 2_000_000,
           Test.test "2 attempts" <| \() ->
             Partition.microSecondsDelayForAttempt 2
-              |> Expect.equal 100_000_000,
+              |> Expect.equal 4_000_000,
           Test.test "3 attempts" <| \() ->
             Partition.microSecondsDelayForAttempt 3
-              |> Expect.equal 1000_000_000,
+              |> Expect.equal 8_000_000,
           Test.test "4 attempts" <| \() ->
             Partition.microSecondsDelayForAttempt 4
-              |> Expect.equal 3_600_000_000
+              |> Expect.equal 16_000_000,
+          Test.test "100 attempts" <| \() ->
+            Partition.microSecondsDelayForAttempt 100
+              |> Expect.equal 1_024_000_000
         ]
     ]
