diff --git a/CHANGELOG.md b/CHANGELOG.md
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -7,6 +7,14 @@
 
 ## Unreleased
 
+## [0.8.0] - 2022-07-10
+
+### Fixed
+
+- Cachix Deploy: properly fix disconnection issues
+- Cachix Deploy: deployments are now a separate process so cachix agent can be upgraded at any time
+
+
 ## [0.7.1] - 2022-06-27
 
 ### Fixed
diff --git a/cachix-deployment/Main.hs b/cachix-deployment/Main.hs
new file mode 100644
--- /dev/null
+++ b/cachix-deployment/Main.hs
@@ -0,0 +1,78 @@
+{-# LANGUAGE DuplicateRecordFields #-}
+
+module Main
+  ( main,
+  )
+where
+
+import Cachix.API.Error (escalateAs)
+import qualified Cachix.API.WebSocketSubprotocol as WSS
+import Cachix.Client.Retry
+import qualified Cachix.Deploy.Activate as Activate
+import qualified Cachix.Deploy.Websocket as CachixWebsocket
+import Conduit ((.|))
+import qualified Control.Concurrent.Async as Async
+import qualified Control.Concurrent.STM.TQueue as TQueue
+import qualified Data.Aeson as Aeson
+import qualified Data.Conduit as Conduit
+import qualified Data.Conduit.Combinators as Conduit
+import qualified Data.Conduit.TQueue as Conduit
+import Data.String (String)
+import Data.Time.Clock (getCurrentTime)
+import Data.UUID (UUID)
+import qualified Data.UUID as UUID
+import GHC.IO.Encoding
+import qualified Katip as K
+import Network.HTTP.Simple (RequestHeaders)
+import qualified Network.WebSockets as WS
+import Protolude hiding (toS)
+import Protolude.Conv
+import System.IO (BufferMode (..), hSetBuffering)
+import qualified Wuss
+
+main :: IO ()
+main = do
+  setLocaleEncoding utf8
+  hSetBuffering stdout LineBuffering
+  hSetBuffering stderr LineBuffering
+  input <- escalateAs (FatalError . toS) . Aeson.eitherDecode . toS =<< getContents
+  CachixWebsocket.runForever (CachixWebsocket.websocketOptions input) (handleMessage input)
+
+handleMessage :: CachixWebsocket.Input -> ByteString -> (K.KatipContextT IO () -> IO ()) -> WS.Connection -> CachixWebsocket.AgentState -> ByteString -> K.KatipContextT IO ()
+handleMessage input payload runKatip connection _ agentToken =
+  CachixWebsocket.parseMessage payload (handleCommand . WSS.command)
+  where
+    deploymentDetails = CachixWebsocket.deploymentDetails input
+    options = CachixWebsocket.websocketOptions input
+    handleCommand :: WSS.BackendCommand -> K.KatipContextT IO ()
+    handleCommand (WSS.Deployment _) =
+      K.logLocM K.ErrorS "cachix-deployment should have never gotten a deployment command directly."
+    handleCommand (WSS.AgentRegistered agentInformation) = do
+      queue <- liftIO $ atomically TQueue.newTQueue
+      let deploymentID = WSS.id (deploymentDetails :: WSS.DeploymentDetails)
+          streamingThread = runLogStreaming (toS $ CachixWebsocket.host options) (CachixWebsocket.headers options agentToken) queue deploymentID
+          activateThread = runKatip $ do
+            Activate.activate options connection (Conduit.sinkTQueue queue) deploymentDetails agentInformation agentToken
+      liftIO $ Async.race_ streamingThread activateThread
+      throwIO ExitSuccess
+    runLogStreaming :: String -> RequestHeaders -> Conduit.TQueue ByteString -> UUID -> IO ()
+    runLogStreaming host headers queue deploymentID = do
+      let path = "/api/v1/deploy/log/" <> UUID.toText deploymentID
+      retryAllWithLogging endlessRetryPolicy (CachixWebsocket.logger runKatip) $ do
+        liftIO $
+          Wuss.runSecureClientWith host 443 (toS path) WS.defaultConnectionOptions headers $
+            \conn ->
+              bracket_ (return ()) (WS.sendClose connection ("Closing." :: ByteString)) $
+                Conduit.runConduit $
+                  Conduit.sourceTQueue queue
+                    .| Conduit.linesUnboundedAscii
+                    -- TODO: prepend katip-like format to each line
+                    -- .| (if CachixWebsocket.isVerbose options then Conduit.print else mempty)
+                    .| sendLog conn
+
+sendLog :: WS.Connection -> Conduit.ConduitT ByteString Conduit.Void IO ()
+sendLog connection = Conduit.mapM_ f
+  where
+    f = \bs -> do
+      now <- getCurrentTime
+      WS.sendTextData connection $ Aeson.encode $ WSS.Log {WSS.line = toS bs, WSS.time = now}
diff --git a/cachix.cabal b/cachix.cabal
--- a/cachix.cabal
+++ b/cachix.cabal
@@ -1,6 +1,6 @@
 cabal-version:      2.2
 name:               cachix
-version: 0.7.1
+version: 0.8.0
 license:            Apache-2.0
 license-file:       LICENSE
 copyright:          2018 Domen Kožar
@@ -66,6 +66,9 @@
     Cachix.Deploy.ActivateCommand
     Cachix.Deploy.Agent
     Cachix.Deploy.OptionsParser
+    Cachix.Deploy.StdinProcess
+    Cachix.Deploy.Websocket
+    Cachix.Deploy.WebsocketPong
     System.Nix.Base32
 
   hs-source-dirs:    src
@@ -145,6 +148,32 @@
     , cachix
     , cachix-api
     , safe-exceptions
+
+executable .cachix-deployment
+  import:             defaults
+  main-is:            Main.hs
+  build-tool-depends: hspec-discover:hspec-discover -any
+  ghc-options:        -threaded -rtsopts -with-rtsopts=-maxN8
+  hs-source-dirs:     cachix-deployment
+  other-modules:      Paths_cachix
+  autogen-modules:    Paths_cachix
+  build-depends:
+    , aeson
+    , async
+    , base             >=4.7 && <5
+    , cachix
+    , cachix-api
+    , conduit
+    , http-conduit
+    , katip
+    , protolude
+    , safe-exceptions
+    , stm
+    , stm-conduit
+    , time
+    , uuid
+    , websockets
+    , wuss
 
 test-suite cachix-test
   import:         defaults
diff --git a/src/Cachix/Client/Commands.hs b/src/Cachix/Client/Commands.hs
--- a/src/Cachix/Client/Commands.hs
+++ b/src/Cachix/Client/Commands.hs
@@ -163,23 +163,26 @@
   pushParams <- getPushParams env opts name
   normalized <-
     liftIO $
-      for inputStorePaths $ \path -> do
-        storePath <- followLinksToStorePath (pushParamsStore pushParams) (encodeUtf8 path)
-        filterInvalidStorePath (pushParamsStore pushParams) storePath
-
+      for inputStorePaths $
+        \path -> do
+          storePath <- followLinksToStorePath (pushParamsStore pushParams) (encodeUtf8 path)
+          filterInvalidStorePath (pushParamsStore pushParams) storePath
   pushedPaths <-
     pushClosure
       (mapConcurrentlyBounded (numJobs opts))
       pushParams
       (catMaybes normalized)
   case (length normalized, length pushedPaths) of
-    (0, _) -> putText "Nothing to push."
-    (_, 0) -> putText "Nothing to push - all store paths are already on Cachix."
-    _ -> putText "All done."
+    (0, _) -> putTextError "Nothing to push."
+    (_, 0) -> putTextError "Nothing to push - all store paths are already on Cachix."
+    _ -> putTextError "All done."
 push _ _ =
   throwIO $
     DeprecatedCommand "DEPRECATED: cachix watch-store has replaced cachix push --watch-store."
 
+putTextError :: Text -> IO ()
+putTextError = hPutStrLn stderr
+
 watchStore :: Env -> PushOptions -> Text -> IO ()
 watchStore env opts name = do
   pushParams <- getPushParams env opts name
@@ -193,7 +196,6 @@
       watch = do
         hDuplicateTo stderr stdout -- redirect all stdout to stderr
         WatchStore.startWorkers (pushParamsStore pushParams) (numJobs pushOpts) pushParams
-
   (_, exitCode) <- concurrently watch $ do
     (_, _, _, processHandle) <- System.Process.createProcess process
     exitCode <- System.Process.waitForProcess processHandle
diff --git a/src/Cachix/Client/Config/Orphans.hs b/src/Cachix/Client/Config/Orphans.hs
--- a/src/Cachix/Client/Config/Orphans.hs
+++ b/src/Cachix/Client/Config/Orphans.hs
@@ -12,7 +12,7 @@
 import Servant.Auth.Client
 
 instance FromDhall Token where
-  autoWith _ = strictText { extract = ex }
+  autoWith _ = strictText {extract = ex}
     where
       ex (TextLit (Chunks [] t)) = pure (Token (toS t))
       ex _ = panic "Unexpected Dhall value. Did it typecheck?"
diff --git a/src/Cachix/Client/Push.hs b/src/Cachix/Client/Push.hs
--- a/src/Cachix/Client/Push.hs
+++ b/src/Cachix/Client/Push.hs
@@ -271,9 +271,10 @@
           )
   let missingHashes = Set.fromList (encodeUtf8 <$> missingHashesList)
   pathsAndHashes <- liftIO $
-    for paths $ \path -> do
-      hash_ <- Store.getStorePathHash path
-      pure (hash_, path)
+    for paths $
+      \path -> do
+        hash_ <- Store.getStorePathHash path
+        pure (hash_, path)
   return $ map snd $ filter (\(hash_, _path) -> Set.member hash_ missingHashes) pathsAndHashes
 
 -- TODO: move to a separate module specific to cli
diff --git a/src/Cachix/Client/PushQueue.hs b/src/Cachix/Client/PushQueue.hs
--- a/src/Cachix/Client/PushQueue.hs
+++ b/src/Cachix/Client/PushQueue.hs
@@ -42,10 +42,11 @@
 worker pushParams workerState = forever $ do
   storePath <- atomically $ TBQueue.readTBQueue $ pushQueue workerState
   bracket_ (inProgresModify (+ 1)) (inProgresModify (\x -> x - 1)) $
-    retryAll $ \retrystatus ->
-      void $ do
-        maybeStorePath <- filterInvalidStorePath (Push.pushParamsStore pushParams) storePath
-        for maybeStorePath $ \validatedStorePath -> Push.uploadStorePath pushParams validatedStorePath retrystatus
+    retryAll $
+      \retrystatus ->
+        void $ do
+          maybeStorePath <- filterInvalidStorePath (Push.pushParamsStore pushParams) storePath
+          for maybeStorePath $ \validatedStorePath -> Push.uploadStorePath pushParams validatedStorePath retrystatus
   where
     inProgresModify f =
       atomically $ modifyTVar' (inProgress workerState) f
@@ -100,7 +101,7 @@
 
 exitOnceQueueIsEmpty :: IO () -> Async () -> Async () -> QueryWorkerState -> PushWorkerState -> IO ()
 exitOnceQueueIsEmpty stopProducerCallback pushWorker queryWorker queryWorkerState pushWorkerState = do
-  putText "Stopped watching /nix/store and waiting for queue to empty ..."
+  putTextError "Stopped watching /nix/store and waiting for queue to empty ..."
   Systemd.notifyStopping
   stopProducerCallback
   go
@@ -115,12 +116,15 @@
         return (isDone, inprogress, pushQueueLength)
       if isDone
         then do
-          putText "Done."
+          putTextError "Done."
           cancelWith queryWorker StopWorker
           cancelWith pushWorker StopWorker
         else do
           -- extend shutdown for another 90s
           Systemd.notify False $ "EXTEND_TIMEOUT_USEC=" <> show (90 * 1000 * 1000)
-          putText $ "Waiting to finish: " <> show inprogress <> " pushing, " <> show queueLength <> " in queue"
+          putTextError $ "Waiting to finish: " <> show inprogress <> " pushing, " <> show queueLength <> " in queue"
           threadDelay (1000 * 1000)
           go
+
+putTextError :: Text -> IO ()
+putTextError = hPutStrLn stderr
diff --git a/src/Cachix/Client/Retry.hs b/src/Cachix/Client/Retry.hs
--- a/src/Cachix/Client/Retry.hs
+++ b/src/Cachix/Client/Retry.hs
@@ -1,3 +1,5 @@
+{-# LANGUAGE ScopedTypeVariables #-}
+
 module Cachix.Client.Retry
   ( retryAll,
     retryAllWithLogging,
@@ -5,9 +7,9 @@
   )
 where
 
-import Control.Exception.Safe (MonadMask, isSyncException)
+import Control.Exception.Safe (Handler (..), MonadMask, isSyncException)
 import Control.Retry (RetryPolicy, RetryPolicyM, RetryStatus, exponentialBackoff, limitRetries, logRetries, recoverAll, recovering, skipAsyncExceptions)
-import Protolude
+import Protolude hiding (Handler (..))
 
 retryAll :: (MonadIO m, MonadMask m) => (RetryStatus -> m a) -> m a
 retryAll =
@@ -17,10 +19,10 @@
 retryAllWithLogging :: (MonadIO m, MonadMask m) => RetryPolicyM m -> (Bool -> SomeException -> RetryStatus -> m ()) -> m a -> m a
 retryAllWithLogging retryPolicy logger action = recovering retryPolicy handlers $ const action
   where
-    handlers = skipAsyncExceptions ++ [loggingHandler]
-
+    handlers = skipAsyncExceptions ++ [exitSuccessHandler, loggingHandler]
+    exitSuccessHandler :: MonadIO m => RetryStatus -> Handler m Bool
+    exitSuccessHandler _ = Handler $ \(_ :: ExitCode) -> return False
     loggingHandler = logRetries exceptionPredicate logger
-
     exceptionPredicate = return . isSyncException
 
 defaultRetryPolicy :: RetryPolicy
diff --git a/src/Cachix/Client/Servant.hs b/src/Cachix/Client/Servant.hs
--- a/src/Cachix/Client/Servant.hs
+++ b/src/Cachix/Client/Servant.hs
@@ -8,20 +8,20 @@
 module Cachix.Client.Servant
   ( isErr,
     cachixClient,
-    deployClient
+    deployClient,
   )
 where
 
 import qualified Cachix.API
+import qualified Cachix.API.Deploy as Cachix.Deploy.API
+import Cachix.Types.ContentTypes ()
 import Network.HTTP.Types (Status)
 import Protolude
 import Servant.API.Generic
 import Servant.Auth.Client ()
-import Cachix.Types.ContentTypes ()
 import qualified Servant.Client
 import Servant.Client.Generic (AsClientT)
 import Servant.Client.Streaming
-import qualified Cachix.API.Deploy as Cachix.Deploy.API
 import Servant.Conduit ()
 
 isErr :: ClientError -> Status -> Bool
diff --git a/src/Cachix/Client/WatchStore.hs b/src/Cachix/Client/WatchStore.hs
--- a/src/Cachix/Client/WatchStore.hs
+++ b/src/Cachix/Client/WatchStore.hs
@@ -20,7 +20,7 @@
 
 producer :: Store -> WatchManager -> PushQueue.Queue -> IO (IO ())
 producer store mgr queue = do
-  putText "Watching /nix/store for new store paths ..."
+  putTextError "Watching /nix/store for new store paths ..."
   watchDir mgr "/nix/store" filterOnlyStorePaths (queueStorePathAction store queue)
 
 queueStorePathAction :: Store -> PushQueue.Queue -> Event -> IO ()
@@ -38,3 +38,6 @@
   | ".drv.lock" `isSuffixOf` fp = False
   | ".lock" `isSuffixOf` fp = True
 filterOnlyStorePaths _ = False
+
+putTextError :: Text -> IO ()
+putTextError = hPutStrLn stderr
diff --git a/src/Cachix/Deploy/Activate.hs b/src/Cachix/Deploy/Activate.hs
--- a/src/Cachix/Deploy/Activate.hs
+++ b/src/Cachix/Deploy/Activate.hs
@@ -5,9 +5,7 @@
 import qualified Cachix.API.WebSocketSubprotocol as WSS
 import qualified Cachix.Client.InstallationMode as InstallationMode
 import qualified Cachix.Client.NetRc as NetRc
-import qualified Cachix.Client.OptionsParser as CachixOptions
-import Cachix.Client.URI (getBaseUrl)
-import qualified Cachix.Deploy.OptionsParser as AgentOptions
+import qualified Cachix.Deploy.Websocket as CachixWebsocket
 import qualified Cachix.Types.BinaryCache as BinaryCache
 import Cachix.Types.Permission (Permission (..))
 import qualified Data.Conduit as Conduit
@@ -20,41 +18,35 @@
 import Protolude hiding (log, toS)
 import Protolude.Conv (toS)
 import Servant.Auth.Client (Token (..))
-import qualified Servant.Client as Servant
 import System.Directory (doesPathExist)
 import System.IO.Temp (withSystemTempDirectory)
 import System.Process
 import Prelude (String)
 
--- TODO: duplicated in Agent.hs
-host :: CachixOptions.CachixOptions -> String
-host cachixOptions = Servant.baseUrlHost $ getBaseUrl $ CachixOptions.host cachixOptions
-
-domain :: CachixOptions.CachixOptions -> WSS.Cache -> Text
-domain cachixOptions cache = toS (WSS.cacheName cache) <> "." <> toS (host cachixOptions)
+domain :: CachixWebsocket.Options -> WSS.Cache -> Text
+domain options cache = toS (WSS.cacheName cache) <> "." <> toS (CachixWebsocket.host options)
 
 -- TODO: get uri scheme
-uri :: CachixOptions.CachixOptions -> WSS.Cache -> Text
-uri cachixOptions cache = "https://" <> domain cachixOptions cache
+uri :: CachixWebsocket.Options -> WSS.Cache -> Text
+uri options cache = "https://" <> domain options cache
 
 -- TODO: what if websocket gets closed while deploying?
 activate ::
-  CachixOptions.CachixOptions ->
-  AgentOptions.AgentOptions ->
+  CachixWebsocket.Options ->
   WS.Connection ->
   Conduit.ConduitT ByteString Void IO () ->
   WSS.DeploymentDetails ->
   WSS.AgentInformation ->
   ByteString ->
   K.KatipContextT IO ()
-activate cachixOptions agentArgs connection sourceStream deploymentDetails agentInfo agentToken = do
+activate options connection sourceStream deploymentDetails agentInfo agentToken = do
   let storePath = WSS.storePath deploymentDetails
       cachesArgs :: [String]
       cachesArgs = case WSS.cache agentInfo of
         Just cache ->
           let officialCache = "cache.nixos.org-1:6NCHdD59X431o0gWypbMrAURkbJ16ZPMQFGspcDShjY="
-              substituters = ["--option", "extra-substituters", toS (uri cachixOptions cache)]
-              sigs = ["--option", "trusted-public-keys", officialCache <> " " <> toS (domain cachixOptions cache) <> "-1:" <> toS (WSS.publicKey cache)]
+              substituters = ["--option", "extra-substituters", toS (uri options cache)]
+              sigs = ["--option", "trusted-public-keys", officialCache <> " " <> toS (domain options cache) <> "-1:" <> toS (WSS.publicKey cache)]
            in substituters ++ sigs
         Nothing -> []
       deploymentID = WSS.id (deploymentDetails :: WSS.DeploymentDetails)
@@ -69,10 +61,8 @@
         K.logLocM K.InfoS $ K.ls $ "Deploying #" <> index <> " failed."
         sendMessage $ deploymentFinished False now
         -- hack to flush logs
-        liftIO $ threadDelay (1 * 1000 * 1000)
-
+        liftIO $ threadDelay (5 * 1000 * 1000)
   K.logLocM K.InfoS $ K.ls $ "Deploying #" <> index <> ": " <> WSS.storePath deploymentDetails
-
   -- notify the service deployment started
   now <- liftIO getCurrentTime
   sendMessage $
@@ -80,36 +70,35 @@
       { WSS.id = deploymentID,
         WSS.time = now
       }
-
   -- TODO: don't create tmpfile for public caches
   -- TODO: add GC root so it's preserved for the next command
   -- get the store path using caches
   (downloadExitCode, _, _) <- liftIO $
-    withSystemTempDirectory "netrc" $ \dir -> do
-      let filepath = dir <> "netrc"
-      args <- case WSS.cache agentInfo of
-        Just cache -> do
-          -- TODO: ugh
-          let bc =
-                BinaryCache.BinaryCache
-                  { BinaryCache.name = "",
-                    BinaryCache.uri = toS (uri cachixOptions cache),
-                    BinaryCache.publicSigningKeys = [],
-                    BinaryCache.isPublic = WSS.isPublic cache,
-                    BinaryCache.githubUsername = "",
-                    BinaryCache.permission = Read
-                  }
-          NetRc.add (Token agentToken) [bc] filepath
-          return $ cachesArgs <> ["--option", "netrc-file", filepath]
-        Nothing ->
-          return cachesArgs
-      shellOut "nix-store" (["-r", toS storePath] <> args)
-
+    withSystemTempDirectory "netrc" $
+      \dir -> do
+        let filepath = dir <> "netrc"
+        args <- case WSS.cache agentInfo of
+          Just cache -> do
+            -- TODO: ugh
+            let bc =
+                  BinaryCache.BinaryCache
+                    { BinaryCache.name = "",
+                      BinaryCache.uri = toS (uri options cache),
+                      BinaryCache.publicSigningKeys = [],
+                      BinaryCache.isPublic = WSS.isPublic cache,
+                      BinaryCache.githubUsername = "",
+                      BinaryCache.permission = Read
+                    }
+            NetRc.add (Token agentToken) [bc] filepath
+            return $ cachesArgs <> ["--option", "netrc-file", filepath]
+          Nothing ->
+            return cachesArgs
+        shellOut "nix-store" (["-r", toS storePath] <> args)
   -- TODO: use exceptions to simplify this code
   case downloadExitCode of
     ExitFailure _ -> deploymentFailed
     ExitSuccess -> do
-      (profile, activationScripts) <- liftIO $ getActivationScript storePath (AgentOptions.profile agentArgs)
+      (profile, activationScripts) <- liftIO $ getActivationScript storePath (CachixWebsocket.profile options)
       -- TODO: document what happens if wrong user is used for the agent
 
       -- set the new profile
@@ -121,18 +110,14 @@
           exitCodes <- for activationScripts $ \(cmd, args) -> do
             (activateScriptExitCode, _, _) <- liftIO $ shellOut cmd args
             return activateScriptExitCode
-
           if not (all (== ExitSuccess) exitCodes)
             then deploymentFailed
             else do
               now <- liftIO getCurrentTime
               sendMessage $ deploymentFinished True now
-
               liftIO $ log "Successfully activated the deployment."
-
               -- TODO: this is a hack to make sure the deployment is finished
-              liftIO $ threadDelay (2 * 1000 * 1000)
-
+              liftIO $ threadDelay (5 * 1000 * 1000)
               K.logLocM K.InfoS $ K.ls $ "Deployment #" <> index <> " finished"
   where
     -- TODO: prevent service from being restarted while deploying
@@ -140,18 +125,14 @@
 
     index :: Text
     index = show $ WSS.index deploymentDetails
-
     shellOut cmd args = do
       log $ "$ " <> toS cmd <> " " <> toS (unwords $ fmap toS args)
       Conduit.sourceProcessWithStreams (proc cmd args) Conduit.sinkNull sourceStream sourceStream
-
     log :: ByteString -> IO ()
     log msg = Conduit.connect (Conduit.yieldMany ["\n" <> msg <> "\n"]) sourceStream
-
     sendMessage cmd = liftIO $ do
       command <- createMessage cmd
       WSS.sendMessage connection command
-
     createMessage command = do
       uuid <- UUID.nextRandom
       return $
@@ -165,7 +146,6 @@
         method = case command of
           WSS.DeploymentStarted {} -> "DeploymentStarted"
           WSS.DeploymentFinished {} -> "DeploymentFinished"
-
     getActivationScript :: Text -> Text -> IO (Text, [Command])
     getActivationScript storePath profile = do
       isNixOS <- doesPathExist $ toS $ storePath <> "/nixos-version"
diff --git a/src/Cachix/Deploy/Agent.hs b/src/Cachix/Deploy/Agent.hs
--- a/src/Cachix/Deploy/Agent.hs
+++ b/src/Cachix/Deploy/Agent.hs
@@ -4,116 +4,52 @@
 
 import qualified Cachix.API.WebSocketSubprotocol as WSS
 import qualified Cachix.Client.OptionsParser as CachixOptions
-import Cachix.Client.Retry
 import Cachix.Client.URI (getBaseUrl)
-import Cachix.Client.Version (versionNumber)
-import qualified Cachix.Deploy.Activate as Activate
 import qualified Cachix.Deploy.OptionsParser as AgentOptions
-import Conduit ((.|))
-import qualified Control.Concurrent.Async as Async
-import qualified Control.Concurrent.STM.TQueue as TQueue
-import qualified Control.Exception.Safe as Exception
+import Cachix.Deploy.StdinProcess (readProcess)
+import qualified Cachix.Deploy.Websocket as CachixWebsocket
 import qualified Data.Aeson as Aeson
-import qualified Data.Conduit as Conduit
-import qualified Data.Conduit.Combinators as Conduit
-import qualified Data.Conduit.TQueue as Conduit
-import Data.IORef
-import Data.String (String)
-import Data.Time.Clock (getCurrentTime)
-import Data.UUID (UUID)
-import qualified Data.UUID as UUID
-import Katip (KatipContextT)
 import qualified Katip as K
-import Network.HTTP.Simple (RequestHeaders)
 import qualified Network.WebSockets as WS
+import Paths_cachix (getBinDir)
 import Protolude hiding (toS)
 import Protolude.Conv
 import qualified Servant.Client as Servant
-import System.Environment (getEnv)
-import qualified Wuss
 
-type AgentState = IORef (Maybe WSS.AgentInformation)
-
 run :: CachixOptions.CachixOptions -> AgentOptions.AgentOptions -> IO ()
-run cachixOptions agentOpts = withKatip (CachixOptions.verbose cachixOptions) $ \logEnv -> do
-  agentToken <- getEnv "CACHIX_AGENT_TOKEN"
-  -- TODO: error if token is missing
-  agentState <- newIORef Nothing
-  let runKatip = K.runKatipContextT logEnv () "agent"
-      headers =
-        [ ("Authorization", "Bearer " <> toS agentToken),
-          ("name", name),
-          ("version", toS versionNumber)
-        ]
-      host = Servant.baseUrlHost $ getBaseUrl $ CachixOptions.host cachixOptions
-      path = "/ws"
-  runKatip $
-    retryAllWithLogging endlessRetryPolicy (logger runKatip) $ do
-      K.logLocM K.InfoS $ K.ls ("Agent " <> agentIdentifier <> " connecting to " <> toS host <> toS path)
-      liftIO $
-        Wuss.runSecureClientWith host 443 path WS.defaultConnectionOptions headers $ \connection -> runKatip $ do
-          K.logLocM K.InfoS "Connected to Cachix Deploy service"
-          liftIO $
-            WS.withPingThread connection 30 (runKatip $ K.logLocM K.DebugS "WebSocket keep-alive ping") $ do
-              WSS.recieveDataConcurrently
-                connection
-                (\message -> Exception.handle (handler runKatip) $ runKatip (handleMessage runKatip host headers message connection agentState (toS agentToken)))
-  where
-    name = toS (AgentOptions.name agentOpts)
-    agentIdentifier = name <> " " <> toS versionNumber
-    logger runKatip _ exception _ = runKatip $ K.logLocM K.ErrorS $ K.ls $ "Retrying due to an exception:" <> displayException exception
-    handleMessage :: (KatipContextT IO () -> IO ()) -> String -> RequestHeaders -> ByteString -> WS.Connection -> AgentState -> ByteString -> KatipContextT IO ()
-    handleMessage runKatip host headers payload connection agentState agentToken = do
-      case WSS.parseMessage payload of
-        (Left err) ->
-          -- TODO: show the bytestring?
-          K.logLocM K.ErrorS $ K.ls $ "Failed to parse websocket payload: " <> err
-        (Right message) ->
-          case WSS.command message of
-            WSS.AgentRegistered agentInformation -> do
-              K.logLocM K.InfoS $ K.ls $ "Agent " <> agentIdentifier <> " registered."
-              liftIO $ atomicWriteIORef agentState (Just agentInformation)
-            WSS.Deployment deploymentDetails -> do
-              maybeAgentInformation <- liftIO $ readIORef agentState
-              let index :: Text
-                  index = show $ WSS.index deploymentDetails
-                  deploymentID = WSS.id (deploymentDetails :: WSS.DeploymentDetails)
-              case maybeAgentInformation of
-                Nothing -> K.logLocM K.InfoS $ K.ls $ "Ignoring deployment #" <> index <> " as agent isn't registered yet."
-                Just agentInformation -> do
-                  queue <- liftIO $ atomically TQueue.newTQueue
-                  liftIO $ Async.race_ (runLogStreaming runKatip host headers queue deploymentID) $ runKatip $ Activate.activate cachixOptions agentOpts connection (Conduit.sinkTQueue queue) deploymentDetails agentInformation agentToken
-
-    runLogStreaming :: (KatipContextT IO () -> IO ()) -> String -> RequestHeaders -> Conduit.TQueue ByteString -> UUID -> IO ()
-    runLogStreaming runKatip host headers queue deploymentID = do
-      -- TODO: debug Conduit.print
-      let path = "/api/v1/deploy/log/" <> UUID.toText deploymentID
-      retryAllWithLogging endlessRetryPolicy (logger runKatip) $ do
-        liftIO $
-          Wuss.runSecureClientWith host 443 (toS path) WS.defaultConnectionOptions headers $ \connection ->
-            bracket_ (return ()) (WS.sendClose connection ("Closing." :: ByteString)) $
-              Conduit.runConduit $
-                Conduit.sourceTQueue queue
-                  .| Conduit.linesUnboundedAscii
-                  .| websocketSend connection
-
-handler :: (KatipContextT IO () -> IO ()) -> Exception.SomeException -> IO ()
-handler runKatip e = do
-  runKatip $ K.logLocM K.ErrorS $ "Unexpected exception: " <> K.ls (Exception.displayException e)
-
-withKatip :: Bool -> (K.LogEnv -> IO a) -> IO a
-withKatip isVerbose =
-  bracket createLogEnv K.closeScribes
-  where
-    permit = if isVerbose then K.DebugS else K.InfoS
-    createLogEnv = do
-      logEnv <- K.initLogEnv "agent" "production"
-      stdoutScribe <- K.mkHandleScribe K.ColorIfTerminal stdout (K.permitItem permit) K.V2
-      K.registerScribe "stdout" stdoutScribe K.defaultScribeSettings logEnv
-
-websocketSend :: WS.Connection -> Conduit.ConduitT ByteString Conduit.Void IO ()
-websocketSend connection = Conduit.mapM_ f
+run cachixOptions agentOpts = do
+  CachixWebsocket.runForever options handleMessage
   where
-    f = \bs -> do
-      now <- getCurrentTime
-      WS.sendTextData connection $ Aeson.encode $ WSS.Log {line = toS bs, time = now}
+    host = toS $ Servant.baseUrlHost $ getBaseUrl $ CachixOptions.host cachixOptions
+    name = AgentOptions.name agentOpts
+    options =
+      CachixWebsocket.Options
+        { CachixWebsocket.host = host,
+          CachixWebsocket.name = name,
+          CachixWebsocket.path = "/ws",
+          CachixWebsocket.profile = AgentOptions.profile agentOpts,
+          CachixWebsocket.isVerbose = CachixOptions.verbose cachixOptions
+        }
+    handleMessage :: ByteString -> (K.KatipContextT IO () -> IO ()) -> WS.Connection -> CachixWebsocket.AgentState -> ByteString -> K.KatipContextT IO ()
+    handleMessage payload _ _ agentState _ = do
+      CachixWebsocket.parseMessage payload (handleCommand . WSS.command)
+      where
+        handleCommand :: WSS.BackendCommand -> K.KatipContextT IO ()
+        handleCommand (WSS.AgentRegistered agentInformation) = do
+          CachixWebsocket.registerAgent agentState agentInformation
+        handleCommand (WSS.Deployment deploymentDetails) = do
+          -- TODO: lock to ensure one deployment at the time
+          let input =
+                CachixWebsocket.Input
+                  { deploymentDetails = deploymentDetails,
+                    websocketOptions =
+                      CachixWebsocket.Options
+                        { host = host,
+                          name = name,
+                          path = "/ws-deployment",
+                          profile = AgentOptions.profile agentOpts,
+                          isVerbose = CachixOptions.verbose cachixOptions
+                        }
+                  }
+          binDir <- toS <$> liftIO getBinDir
+          liftIO $ readProcess (binDir <> "/.cachix-deployment") [] (toS $ Aeson.encode input)
diff --git a/src/Cachix/Deploy/StdinProcess.hs b/src/Cachix/Deploy/StdinProcess.hs
new file mode 100644
--- /dev/null
+++ b/src/Cachix/Deploy/StdinProcess.hs
@@ -0,0 +1,25 @@
+module Cachix.Deploy.StdinProcess where
+
+import Protolude hiding (stdin)
+import System.IO (hClose)
+import System.Process
+import Prelude (String, userError)
+
+-- | Run a process with only stdin as an input
+readProcess :: FilePath -> [String] -> String -> IO ()
+readProcess cmd args input = do
+  (Just stdin, _, _, ph) <-
+    createProcess
+      (proc cmd args)
+        { std_in = CreatePipe,
+          -- When launching cachix-deployment we need to make sure it's not killed when the main process is killed
+          create_group = True,
+          -- Same, but with posix api
+          new_session = True
+        }
+  hPutStr stdin input
+  hClose stdin
+  exitcode <- waitForProcess ph
+  case exitcode of
+    ExitSuccess -> return ()
+    ExitFailure code -> ioError $ userError $ "Process " <> show cmd <> " failed with exit code " <> show code
diff --git a/src/Cachix/Deploy/Websocket.hs b/src/Cachix/Deploy/Websocket.hs
new file mode 100644
--- /dev/null
+++ b/src/Cachix/Deploy/Websocket.hs
@@ -0,0 +1,103 @@
+-- high level interface for websocket clients
+module Cachix.Deploy.Websocket where
+
+import Cachix.API.WebSocketSubprotocol (AgentInformation)
+import qualified Cachix.API.WebSocketSubprotocol as WSS
+import Cachix.Client.Retry
+import Cachix.Client.Version (versionNumber)
+import qualified Cachix.Deploy.WebsocketPong as WebsocketPong
+import Data.Aeson (FromJSON, ToJSON)
+import Data.IORef
+import qualified Katip as K
+import Network.HTTP.Types (Header)
+import qualified Network.WebSockets as WS
+import Protolude hiding (toS)
+import Protolude.Conv
+import System.Environment (getEnv)
+import qualified Wuss
+
+type AgentState = IORef (Maybe WSS.AgentInformation)
+
+data Options = Options
+  { host :: Text,
+    path :: Text,
+    name :: Text,
+    isVerbose :: Bool,
+    profile :: Text
+  }
+  deriving (Show, Generic, ToJSON, FromJSON)
+
+data Input = Input
+  { deploymentDetails :: WSS.DeploymentDetails,
+    websocketOptions :: Options
+  }
+  deriving (Show, Generic, ToJSON, FromJSON)
+
+runForever :: Options -> (ByteString -> (K.KatipContextT IO () -> IO ()) -> WS.Connection -> AgentState -> ByteString -> K.KatipContextT IO ()) -> IO ()
+runForever options cmd = withKatip (isVerbose options) $ \logEnv -> do
+  agentToken <- getEnv "CACHIX_AGENT_TOKEN"
+  -- TODO: error if token is missing
+  agentState <- newIORef Nothing
+  pongState <- WebsocketPong.newState
+  mainThreadID <- myThreadId
+  let runKatip = K.runKatipContextT logEnv () "agent"
+      pingHandler = do
+        last <- WebsocketPong.secondsSinceLastPong pongState
+        runKatip $ K.logLocM K.DebugS $ K.ls $ "Sending WebSocket keep-alive ping, last pong was " <> (show last :: Text) <> " seconds ago"
+        WebsocketPong.pingHandler pongState mainThreadID pongTimeout
+      connectionOptions = WebsocketPong.installPongHandler pongState WS.defaultConnectionOptions
+  runKatip $
+    retryAllWithLogging endlessRetryPolicy (logger runKatip) $
+      do
+        K.logLocM K.InfoS $ K.ls ("Agent " <> agentIdentifier <> " connecting to " <> toS (host options) <> toS (path options))
+        liftIO $ do
+          -- refresh pong state in case we're reconnecting
+          WebsocketPong.pongHandler pongState
+          Wuss.runSecureClientWith (toS $ host options) 443 (toS $ path options) connectionOptions (headers options (toS agentToken)) $ \connection -> runKatip $ do
+            K.logLocM K.InfoS "Connected to Cachix Deploy service"
+            liftIO $
+              WS.withPingThread connection pingEvery pingHandler $
+                do
+                  WSS.recieveDataConcurrently
+                    connection
+                    (\message -> runKatip (cmd message runKatip connection agentState (toS agentToken)))
+  where
+    agentIdentifier = name options <> " " <> toS versionNumber
+    pingEvery = 30
+    pongTimeout = pingEvery * 2
+
+headers :: Options -> ByteString -> [Header]
+headers options agentToken =
+  [ ("Authorization", "Bearer " <> toS agentToken),
+    ("name", toS (name options)),
+    ("version", toS versionNumber)
+  ]
+
+-- TODO: log the exception
+logger runKatip _ exception _ = runKatip $ K.logLocM K.ErrorS $ K.ls $ "Retrying due to an exception:" <> displayException exception
+
+withKatip :: Bool -> (K.LogEnv -> IO a) -> IO a
+withKatip isVerbose =
+  bracket createLogEnv K.closeScribes
+  where
+    permit = if isVerbose then K.DebugS else K.InfoS
+    createLogEnv = do
+      logEnv <- K.initLogEnv "agent" "production"
+      stdoutScribe <- K.mkHandleScribe K.ColorIfTerminal stdout (K.permitItem permit) K.V2
+      K.registerScribe "stdout" stdoutScribe K.defaultScribeSettings logEnv
+
+parseMessage :: FromJSON cmd => ByteString -> (WSS.Message cmd -> K.KatipContextT IO ()) -> K.KatipContextT IO ()
+parseMessage payload m = do
+  case WSS.parseMessage payload of
+    (Left err) ->
+      -- TODO: show the bytestring?
+      K.logLocM K.ErrorS $ K.ls $ "Failed to parse websocket payload: " <> err
+    (Right message) ->
+      m message
+
+-- commands
+
+registerAgent :: AgentState -> AgentInformation -> K.KatipContextT IO ()
+registerAgent agentState agentInformation = do
+  K.logLocM K.InfoS "Agent registered."
+  liftIO $ atomicWriteIORef agentState (Just agentInformation)
diff --git a/src/Cachix/Deploy/WebsocketPong.hs b/src/Cachix/Deploy/WebsocketPong.hs
new file mode 100644
--- /dev/null
+++ b/src/Cachix/Deploy/WebsocketPong.hs
@@ -0,0 +1,42 @@
+-- Implement ppong on the client side for WS
+-- TODO: upstream to https://github.com/jaspervdj/websockets/issues/159
+module Cachix.Deploy.WebsocketPong where
+
+import Data.IORef
+import Data.Time.Clock (UTCTime, diffUTCTime, getCurrentTime, nominalDiffTimeToSeconds)
+import qualified Network.WebSockets as WS
+import Protolude
+
+type LastPongState = IORef UTCTime
+
+data WebsocketPongTimeout
+  = WebsocketPongTimeout
+  deriving (Show)
+
+instance Exception WebsocketPongTimeout
+
+newState :: IO LastPongState
+newState = do
+  now <- getCurrentTime
+  newIORef now
+
+-- everytime we send a ping we check if we also got a pong back
+pingHandler :: LastPongState -> ThreadId -> Int -> IO ()
+pingHandler state threadID maxLastPing = do
+  last <- secondsSinceLastPong state
+  when (last > maxLastPing) $ do
+    throwTo threadID WebsocketPongTimeout
+
+secondsSinceLastPong :: LastPongState -> IO Int
+secondsSinceLastPong state = do
+  now <- getCurrentTime
+  last <- readIORef state
+  return $ ceiling $ nominalDiffTimeToSeconds $ diffUTCTime now last
+
+pongHandler :: LastPongState -> IO ()
+pongHandler state = do
+  now <- getCurrentTime
+  writeIORef state now
+
+installPongHandler :: LastPongState -> WS.ConnectionOptions -> WS.ConnectionOptions
+installPongHandler state opts = opts {WS.connectionOnPong = pongHandler state}
