packages feed

natskell-1.2.0.0: jetstream/JetStream/Stream/Types.hs

{-# LANGUAGE OverloadedStrings #-}

module JetStream.Stream.Types
  ( StreamAPI (..)
  , Stream (..)
  , RetentionPolicy (..)
  , StorageType (..)
  , DiscardPolicy (..)
  , StreamCompression (..)
  , StreamConfig (..)
  , StreamConfigOption
  , StreamConfigRequest
  , streamConfigRequest
  , validateStreamConfigRequest
  , withRetention
  , withStorage
  , withDiscard
  , withDescription
  , withMaxConsumers
  , withMaxMessages
  , withMaxMessagesPerSubject
  , withMaxBytes
  , withMaxAge
  , withMaxMessageSize
  , withReplicas
  , withDuplicateWindow
  , withDenyDelete
  , withAllowRollup
  , withAllowDirect
  , withCompression
  , PurgeStreamOption
  , purgeStreamRequest
  , withPurgeSubject
  , withPurgeSequence
  , withPurgeKeep
  , StreamListOption
  , streamListRequest
  , streamNamesRequest
  , withStreamListOffset
  , withStreamListSubject
  , StreamMessage (..)
  , StreamMessageSelector (..)
  , streamMessageGetRequest
  , StreamMessageDeleteMode (..)
  , streamMessageDeleteRequest
  , DeleteStreamMessageResponse (..)
  , StreamInfo (..)
  , StreamState (..)
  , StreamCluster (..)
  , StreamPeer (..)
  , StreamSourceInfo (..)
  , DeleteStreamResponse (..)
  , PurgeStreamResponse (..)
  , StreamListResponse (..)
  , StreamNamesResponse (..)
  , durationToNanoseconds
  , nanosecondsToDuration
  ) where

import           Data.Aeson
import           Data.Aeson.Types       (Pair, Parser)
import qualified Data.ByteString        as BS
import qualified Data.ByteString.Base64 as Base64
import           Data.Int               (Int32)
import           Data.Maybe             (catMaybes)
import qualified Data.Text              as T
import           Data.Time.Clock        (NominalDiffTime, UTCTime)
import           Data.Word              (Word64)
import           JetStream.Error        (JetStreamError (JetStreamDecodeError))
import           JetStream.Types
    ( CallOption
    , DiscardPolicy (..)
    , JetStreamRequestOption
    , Payload
    , RetentionPolicy (..)
    , StorageType (..)
    , StreamName
    , Subject
    , applyCallOptions
    , byteStringToJSON
    , diffTimeToNanoseconds
    , parseByteString
    )

-- | Stream management operations. The public module exposes this type
-- abstractly so operations can be added without changing its constructor.
data StreamAPI = StreamAPI
                   { create :: StreamName -> [Subject] -> [StreamConfigOption] -> [JetStreamRequestOption] -> IO (Either JetStreamError StreamInfo)
                   , createOrUpdate :: StreamName -> [Subject] -> [StreamConfigOption] -> [JetStreamRequestOption] -> IO (Either JetStreamError StreamInfo)
                   , update :: StreamName -> [Subject] -> [StreamConfigOption] -> [JetStreamRequestOption] -> IO (Either JetStreamError StreamInfo)
                   , info :: StreamName -> [JetStreamRequestOption] -> IO (Either JetStreamError StreamInfo)
                   , getMessage :: StreamName -> StreamMessageSelector -> [JetStreamRequestOption] -> IO (Either JetStreamError StreamMessage)
                   , deleteMessage :: StreamName -> Word64 -> StreamMessageDeleteMode -> [JetStreamRequestOption] -> IO (Either JetStreamError DeleteStreamMessageResponse)
                   , delete :: StreamName -> [JetStreamRequestOption] -> IO (Either JetStreamError DeleteStreamResponse)
                   , purge :: StreamName -> [PurgeStreamOption] -> [JetStreamRequestOption] -> IO (Either JetStreamError PurgeStreamResponse)
                   , list :: [StreamListOption] -> [JetStreamRequestOption] -> IO (Either JetStreamError StreamListResponse)
                   , names :: [StreamListOption] -> [JetStreamRequestOption] -> IO (Either JetStreamError StreamNamesResponse)
                   }

-- | A stable reference to a stream. It is intentionally just an identity;
-- operations remain on 'StreamAPI', so handles do not retain resources.
newtype Stream = Stream { streamName :: StreamName }
  deriving (Eq, Ord, Show)

data StreamConfigRequest = StreamConfigRequest
                             { streamConfigRequestName :: StreamName
                             , streamConfigRequestSubjects :: [Subject]
                             , streamConfigRequestDescription :: Maybe BS.ByteString
                             , streamConfigRequestRetention :: Maybe RetentionPolicy
                             , streamConfigRequestStorage :: Maybe StorageType
                             , streamConfigRequestDiscard :: Maybe DiscardPolicy
                             , streamConfigRequestMaxConsumers :: Maybe Int
                             , streamConfigRequestMaxMessages :: Maybe Integer
                             , streamConfigRequestMaxMessagesPerSubject :: Maybe Integer
                             , streamConfigRequestMaxBytes :: Maybe Integer
                             , streamConfigRequestMaxAge :: Maybe NominalDiffTime
                             , streamConfigRequestMaxMessageSize :: Maybe Int32
                             , streamConfigRequestReplicas :: Maybe Int
                             , streamConfigRequestDuplicateWindow :: Maybe NominalDiffTime
                             , streamConfigRequestDenyDelete :: Maybe Bool
                             , streamConfigRequestAllowRollup :: Maybe Bool
                             , streamConfigRequestAllowDirect :: Maybe Bool
                             , streamConfigRequestCompression :: Maybe StreamCompression
                             }
  deriving (Eq, Show)

type StreamConfigOption = CallOption StreamConfigRequest

streamConfigRequest :: StreamName -> [Subject] -> [StreamConfigOption] -> StreamConfigRequest
streamConfigRequest name subjects options =
  applyCallOptions options $
    StreamConfigRequest
      { streamConfigRequestName = name
      , streamConfigRequestSubjects = subjects
      , streamConfigRequestDescription = Nothing
      , streamConfigRequestRetention = Nothing
      , streamConfigRequestStorage = Nothing
      , streamConfigRequestDiscard = Nothing
      , streamConfigRequestMaxConsumers = Nothing
      , streamConfigRequestMaxMessages = Nothing
      , streamConfigRequestMaxMessagesPerSubject = Nothing
      , streamConfigRequestMaxBytes = Nothing
      , streamConfigRequestMaxAge = Nothing
      , streamConfigRequestMaxMessageSize = Nothing
      , streamConfigRequestReplicas = Nothing
      , streamConfigRequestDuplicateWindow = Nothing
      , streamConfigRequestDenyDelete = Nothing
      , streamConfigRequestAllowRollup = Nothing
      , streamConfigRequestAllowDirect = Nothing
      , streamConfigRequestCompression = Nothing
      }

validateStreamConfigRequest :: StreamConfigRequest -> Either JetStreamError ()
validateStreamConfigRequest config =
  validateLowerBound "stream max consumers" (streamConfigRequestMaxConsumers config) >>
    validateLowerBound "stream max messages" (streamConfigRequestMaxMessages config) >>
    validateLowerBound "stream max messages per subject" (streamConfigRequestMaxMessagesPerSubject config) >>
    validateLowerBound "stream max bytes" (streamConfigRequestMaxBytes config) >>
    case streamConfigRequestMaxMessageSize config of
      Just maxMessageSize
        | maxMessageSize < (-1) ->
            Left (JetStreamDecodeError "stream max message size must be -1 or greater")
      _ ->
        Right ()
  where
    validateLowerBound label value =
      case value of
        Just number
          | number < (-1) ->
              Left (JetStreamDecodeError (label ++ " must be -1 or greater"))
        _ ->
          Right ()

withRetention :: RetentionPolicy -> StreamConfigOption
withRetention retention config =
  config { streamConfigRequestRetention = Just retention }

withStorage :: StorageType -> StreamConfigOption
withStorage storage config =
  config { streamConfigRequestStorage = Just storage }

withDiscard :: DiscardPolicy -> StreamConfigOption
withDiscard discard config =
  config { streamConfigRequestDiscard = Just discard }

withDescription :: BS.ByteString -> StreamConfigOption
withDescription description config =
  config { streamConfigRequestDescription = Just description }

withMaxConsumers :: Int -> StreamConfigOption
withMaxConsumers maxConsumers config =
  config { streamConfigRequestMaxConsumers = Just maxConsumers }

withMaxMessages :: Integer -> StreamConfigOption
withMaxMessages maxMessages config =
  config { streamConfigRequestMaxMessages = Just maxMessages }

withMaxMessagesPerSubject :: Integer -> StreamConfigOption
withMaxMessagesPerSubject maxMessages config =
  config { streamConfigRequestMaxMessagesPerSubject = Just maxMessages }

withMaxBytes :: Integer -> StreamConfigOption
withMaxBytes maxBytes config =
  config { streamConfigRequestMaxBytes = Just maxBytes }

withMaxAge :: NominalDiffTime -> StreamConfigOption
withMaxAge maxAge config =
  config { streamConfigRequestMaxAge = Just maxAge }

-- | Limit the total encoded bytes of one stored message, including headers.
-- Use @-1@ for unlimited.
withMaxMessageSize :: Int32 -> StreamConfigOption
withMaxMessageSize maxMessageSize config =
  config { streamConfigRequestMaxMessageSize = Just maxMessageSize }

withReplicas :: Int -> StreamConfigOption
withReplicas replicas config =
  config { streamConfigRequestReplicas = Just replicas }

withDuplicateWindow :: NominalDiffTime -> StreamConfigOption
withDuplicateWindow window config =
  config { streamConfigRequestDuplicateWindow = Just window }

withDenyDelete :: Bool -> StreamConfigOption
withDenyDelete denyDelete config =
  config { streamConfigRequestDenyDelete = Just denyDelete }

withAllowRollup :: Bool -> StreamConfigOption
withAllowRollup allowRollup config =
  config { streamConfigRequestAllowRollup = Just allowRollup }

withAllowDirect :: Bool -> StreamConfigOption
withAllowDirect allowDirect config =
  config { streamConfigRequestAllowDirect = Just allowDirect }

withCompression :: StreamCompression -> StreamConfigOption
withCompression compression config =
  config { streamConfigRequestCompression = Just compression }

data StreamCompression = NoCompression
                       | S2Compression
                       | StreamCompressionUnknown T.Text
  deriving (Eq, Show)

data StreamConfig = StreamConfig
                      { streamConfigName :: StreamName
                      , streamConfigSubjects :: Maybe [Subject]
                      , streamConfigDescription :: Maybe BS.ByteString
                      , streamConfigRetention :: RetentionPolicy
                      , streamConfigStorage :: StorageType
                      , streamConfigDiscard :: DiscardPolicy
                      , streamConfigMaxConsumers :: Int
                      , streamConfigMaxMessages :: Integer
                      , streamConfigMaxMessagesPerSubject :: Integer
                      , streamConfigMaxBytes :: Integer
                      , streamConfigMaxAge :: NominalDiffTime
                      , streamConfigMaxMessageSize :: Int32
                      , streamConfigReplicas :: Int
                      , streamConfigDuplicateWindow :: Maybe NominalDiffTime
                      , streamConfigDenyDelete :: Bool
                      , streamConfigAllowRollup :: Bool
                      , streamConfigAllowDirect :: Bool
                      , streamConfigCompression :: StreamCompression
                      }
  deriving (Eq, Show)

data PurgeStreamRequest = PurgeStreamRequest
                            { purgeStreamSubject  :: Maybe Subject
                            , purgeStreamSequence :: Maybe Word64
                            , purgeStreamKeep     :: Maybe Integer
                            }
  deriving (Eq, Show)

type PurgeStreamOption = CallOption PurgeStreamRequest

purgeStreamRequest :: [PurgeStreamOption] -> PurgeStreamRequest
purgeStreamRequest options =
  applyCallOptions options $
    PurgeStreamRequest
      { purgeStreamSubject = Nothing
      , purgeStreamSequence = Nothing
      , purgeStreamKeep = Nothing
      }

withPurgeSubject :: Subject -> PurgeStreamOption
withPurgeSubject subject request =
  request { purgeStreamSubject = Just subject }

withPurgeSequence :: Word64 -> PurgeStreamOption
withPurgeSequence sequenceNumber request =
  request { purgeStreamSequence = Just sequenceNumber }

withPurgeKeep :: Integer -> PurgeStreamOption
withPurgeKeep keep request =
  request { purgeStreamKeep = Just keep }

data StreamListRequest = StreamListRequest
                           { streamListRequestOffset  :: Maybe Int
                           , streamListRequestSubject :: Maybe Subject
                           }
  deriving (Eq, Show)

type StreamListOption = CallOption StreamListRequest

streamListRequest :: [StreamListOption] -> StreamListRequest
streamListRequest options =
  applyCallOptions options defaultStreamListRequest

streamNamesRequest :: [StreamListOption] -> StreamListRequest
streamNamesRequest =
  streamListRequest

defaultStreamListRequest :: StreamListRequest
defaultStreamListRequest =
  StreamListRequest
    { streamListRequestOffset = Nothing
    , streamListRequestSubject = Nothing
    }

withStreamListOffset :: Int -> StreamListOption
withStreamListOffset offset request =
  request { streamListRequestOffset = Just offset }

withStreamListSubject :: Subject -> StreamListOption
withStreamListSubject subject request =
  request { streamListRequestSubject = Just subject }

data StreamInfo = StreamInfo
                    { streamInfoConfig  :: StreamConfig
                    , streamInfoCreated :: UTCTime
                    , streamInfoState   :: StreamState
                    , streamInfoCluster :: Maybe StreamCluster
                    , streamInfoMirror  :: Maybe StreamSourceInfo
                    , streamInfoSources :: [StreamSourceInfo]
                    }
  deriving (Eq, Show)

data StreamState = StreamState
                     { streamStateMessages      :: Integer
                     , streamStateBytes         :: Integer
                     , streamStateFirstSequence :: Word64
                     , streamStateFirstTime     :: UTCTime
                     , streamStateLastSequence  :: Word64
                     , streamStateLastTime      :: UTCTime
                     , streamStateConsumerCount :: Int
                     , streamStateDeleted       :: [Word64]
                     , streamStateNumDeleted    :: Integer
                     , streamStateNumSubjects   :: Integer
                     }
  deriving (Eq, Show)

data StreamCluster = StreamCluster
                       { streamClusterName     :: BS.ByteString
                       , streamClusterLeader   :: BS.ByteString
                       , streamClusterReplicas :: [StreamPeer]
                       }
  deriving (Eq, Show)

data StreamPeer = StreamPeer
                    { streamPeerName    :: BS.ByteString
                    , streamPeerCurrent :: Bool
                    , streamPeerOffline :: Bool
                    , streamPeerActive  :: NominalDiffTime
                    , streamPeerLag     :: Integer
                    }
  deriving (Eq, Show)

data StreamSourceInfo = StreamSourceInfo
                          { streamSourceInfoName          :: StreamName
                          , streamSourceInfoFilterSubject :: Subject
                          , streamSourceInfoLag           :: Integer
                          , streamSourceInfoActive        :: NominalDiffTime
                          }
  deriving (Eq, Show)

newtype DeleteStreamResponse = DeleteStreamResponse { deleteStreamSuccess :: Bool }
  deriving (Eq, Show)

data PurgeStreamResponse = PurgeStreamResponse
                             { purgeStreamSuccess :: Bool
                             , purgeStreamPurged  :: Integer
                             }
  deriving (Eq, Show)

data StreamListResponse = StreamListResponse
                            { streamListTotal   :: Int
                            , streamListOffset  :: Int
                            , streamListLimit   :: Int
                            , streamListStreams :: [StreamInfo]
                            }
  deriving (Eq, Show)

data StreamNamesResponse = StreamNamesResponse
                             { streamNamesTotal   :: Int
                             , streamNamesOffset  :: Int
                             , streamNamesLimit   :: Int
                             , streamNamesStreams :: [StreamName]
                             }
  deriving (Eq, Show)

data StreamMessage = StreamMessage
                       { streamMessageSubject    :: Subject
                       , streamMessageSequence   :: Word64
                       , streamMessageHeadersRaw :: Maybe BS.ByteString
                       , streamMessagePayload    :: Maybe Payload
                       , streamMessageTime       :: UTCTime
                       }
  deriving (Eq, Show)

data StreamMessageSelector = StreamMessageBySequence Word64
                           | LastStreamMessageForSubject Subject
                           | NextStreamMessageForSubject Subject
  deriving (Eq, Show)

newtype StreamMessageGetRequest = StreamMessageGetRequest { streamMessageGetSelector :: StreamMessageSelector }
  deriving (Eq, Show)

streamMessageGetRequest :: StreamMessageSelector -> StreamMessageGetRequest
streamMessageGetRequest =
  StreamMessageGetRequest

data StreamMessageDeleteMode = DeleteMessage | SecureDeleteMessage
  deriving (Eq, Show)

data StreamMessageDeleteRequest = StreamMessageDeleteRequest
                                    { streamMessageDeleteSequence :: Word64
                                    , streamMessageDeleteNoErase  :: Maybe Bool
                                    }
  deriving (Eq, Show)

streamMessageDeleteRequest :: Word64 -> StreamMessageDeleteMode -> StreamMessageDeleteRequest
streamMessageDeleteRequest sequenceNumber mode =
  StreamMessageDeleteRequest
    { streamMessageDeleteSequence = sequenceNumber
    , streamMessageDeleteNoErase =
        case mode of
          DeleteMessage       -> Just True
          SecureDeleteMessage -> Nothing
    }

newtype DeleteStreamMessageResponse = DeleteStreamMessageResponse { deleteStreamMessageSuccess :: Bool }
  deriving (Eq, Show)

instance ToJSON StreamConfigRequest where
  toJSON config =
    object $
      [ byteStringPair "name" (streamConfigRequestName config)
      , byteStringListPair "subjects" (streamConfigRequestSubjects config)
      ] ++ catMaybes
        [ maybeByteStringPair "description" (streamConfigRequestDescription config)
        , maybePair "retention" (streamConfigRequestRetention config)
        , maybePair "storage" (streamConfigRequestStorage config)
        , maybePair "discard" (streamConfigRequestDiscard config)
        , maybePair "max_consumers" (streamConfigRequestMaxConsumers config)
        , maybePair "max_msgs" (streamConfigRequestMaxMessages config)
        , maybePair "max_msgs_per_subject" (streamConfigRequestMaxMessagesPerSubject config)
        , maybePair "max_bytes" (streamConfigRequestMaxBytes config)
        , maybeDurationPair "max_age" (streamConfigRequestMaxAge config)
        , maybePair "max_msg_size" (streamConfigRequestMaxMessageSize config)
        , maybePair "num_replicas" (streamConfigRequestReplicas config)
        , maybeDurationPair "duplicate_window" (streamConfigRequestDuplicateWindow config)
        , maybePair "deny_delete" (streamConfigRequestDenyDelete config)
        , maybePair "allow_rollup_hdrs" (streamConfigRequestAllowRollup config)
        , maybePair "allow_direct" (streamConfigRequestAllowDirect config)
        , maybePair "compression" (streamConfigRequestCompression config)
        ]

instance ToJSON StreamConfig where
  toJSON config =
    object . catMaybes $
      [ Just (byteStringPair "name" (streamConfigName config))
      , maybeByteStringListPair "subjects" (streamConfigSubjects config)
      , maybeByteStringPair "description" (streamConfigDescription config)
      , Just ("retention" .= streamConfigRetention config)
      , Just ("storage" .= streamConfigStorage config)
      , Just ("discard" .= streamConfigDiscard config)
      , Just ("max_consumers" .= streamConfigMaxConsumers config)
      , Just ("max_msgs" .= streamConfigMaxMessages config)
      , Just ("max_msgs_per_subject" .= streamConfigMaxMessagesPerSubject config)
      , Just ("max_bytes" .= streamConfigMaxBytes config)
      , Just ("max_age" .= durationToNanoseconds (streamConfigMaxAge config))
      , Just ("max_msg_size" .= streamConfigMaxMessageSize config)
      , Just ("num_replicas" .= streamConfigReplicas config)
      , maybeDurationPair "duplicate_window" (streamConfigDuplicateWindow config)
      , Just ("deny_delete" .= streamConfigDenyDelete config)
      , Just ("allow_rollup_hdrs" .= streamConfigAllowRollup config)
      , Just ("allow_direct" .= streamConfigAllowDirect config)
      , Just ("compression" .= streamConfigCompression config)
      ]

instance FromJSON StreamConfig where
  parseJSON =
    withObject "StreamConfig" $ \value ->
      StreamConfig
        <$> parseByteStringField value "name"
        <*> parseOptionalByteStringListField value "subjects"
        <*> parseOptionalByteStringField value "description"
        <*> value .: "retention"
        <*> value .: "storage"
        <*> value .: "discard"
        <*> value .:? "max_consumers" .!= (-1)
        <*> value .: "max_msgs"
        <*> value .:? "max_msgs_per_subject" .!= (-1)
        <*> value .: "max_bytes"
        <*> parseDurationField value "max_age"
        <*> value .:? "max_msg_size" .!= (-1)
        <*> value .: "num_replicas"
        <*> parseOptionalDurationField value "duplicate_window"
        <*> value .:? "deny_delete" .!= False
        <*> value .:? "allow_rollup_hdrs" .!= False
        <*> value .: "allow_direct"
        <*> value .:? "compression" .!= NoCompression

instance ToJSON StreamCompression where
  toJSON compression =
    String $
      case compression of
        NoCompression                 -> "none"
        S2Compression                 -> "s2"
        StreamCompressionUnknown name -> name

instance FromJSON StreamCompression where
  parseJSON = withText "StreamCompression" $ \value ->
    case value of
      "none" -> pure NoCompression
      "s2"   -> pure S2Compression
      _      -> pure (StreamCompressionUnknown value)

instance ToJSON PurgeStreamRequest where
  toJSON request =
    object . catMaybes $
      [ maybeByteStringPair "filter" (purgeStreamSubject request)
      , maybePair "seq" (purgeStreamSequence request)
      , maybePair "keep" (purgeStreamKeep request)
      ]

instance ToJSON StreamListRequest where
  toJSON request =
    object . catMaybes $
      [ maybePair "offset" (streamListRequestOffset request)
      , maybeByteStringPair "subject" (streamListRequestSubject request)
      ]

instance ToJSON StreamMessageGetRequest where
  toJSON request =
    case streamMessageGetSelector request of
      StreamMessageBySequence sequenceNumber ->
        object ["seq" .= sequenceNumber]
      LastStreamMessageForSubject subject ->
        object [byteStringPair "last_by_subj" subject]
      NextStreamMessageForSubject subject ->
        object [byteStringPair "next_by_subj" subject]

instance ToJSON StreamMessageDeleteRequest where
  toJSON request =
    object . catMaybes $
      [ Just ("seq" .= streamMessageDeleteSequence request)
      , maybePair "no_erase" (streamMessageDeleteNoErase request)
      ]

instance ToJSON StreamInfo where
  toJSON info =
    object . catMaybes $
      [ Just ("config" .= streamInfoConfig info)
      , Just ("created" .= streamInfoCreated info)
      , Just ("state" .= streamInfoState info)
      , maybePair "cluster" (streamInfoCluster info)
      , maybePair "mirror" (streamInfoMirror info)
      , Just ("sources" .= streamInfoSources info)
      ]

instance FromJSON StreamInfo where
  parseJSON =
    withObject "StreamInfo" $ \value ->
      StreamInfo
        <$> value .: "config"
        <*> value .: "created"
        <*> value .: "state"
        <*> value .:? "cluster"
        <*> value .:? "mirror"
        <*> value .:? "sources" .!= []

instance ToJSON StreamState where
  toJSON state =
    object
      [ "messages" .= streamStateMessages state
      , "bytes" .= streamStateBytes state
      , "first_seq" .= streamStateFirstSequence state
      , "first_ts" .= streamStateFirstTime state
      , "last_seq" .= streamStateLastSequence state
      , "last_ts" .= streamStateLastTime state
      , "consumer_count" .= streamStateConsumerCount state
      , "deleted" .= streamStateDeleted state
      , "num_deleted" .= streamStateNumDeleted state
      , "num_subjects" .= streamStateNumSubjects state
      ]

instance FromJSON StreamState where
  parseJSON =
    withObject "StreamState" $ \value ->
      StreamState
        <$> value .: "messages"
        <*> value .: "bytes"
        <*> value .: "first_seq"
        <*> value .: "first_ts"
        <*> value .: "last_seq"
        <*> value .: "last_ts"
        <*> value .: "consumer_count"
        <*> value .:? "deleted" .!= []
        <*> value .:? "num_deleted" .!= 0
        <*> value .:? "num_subjects" .!= 0

instance ToJSON StreamCluster where
  toJSON cluster =
    object
      [ "name" .= byteStringToJSON (streamClusterName cluster)
      , "leader" .= byteStringToJSON (streamClusterLeader cluster)
      , "replicas" .= streamClusterReplicas cluster
      ]

instance FromJSON StreamCluster where
  parseJSON =
    withObject "StreamCluster" $ \value ->
      StreamCluster
        <$> parseOptionalByteStringField value "name" .!= ""
        <*> parseOptionalByteStringField value "leader" .!= ""
        <*> value .:? "replicas" .!= []

instance ToJSON StreamPeer where
  toJSON peer =
    object
      [ "name" .= byteStringToJSON (streamPeerName peer)
      , "current" .= streamPeerCurrent peer
      , "offline" .= streamPeerOffline peer
      , "active" .= durationToNanoseconds (streamPeerActive peer)
      , "lag" .= streamPeerLag peer
      ]

instance FromJSON StreamPeer where
  parseJSON =
    withObject "StreamPeer" $ \value ->
      StreamPeer
        <$> parseByteStringField value "name"
        <*> value .: "current"
        <*> value .:? "offline" .!= False
        <*> parseDurationField value "active"
        <*> value .:? "lag" .!= 0

instance ToJSON StreamSourceInfo where
  toJSON source =
    object
      [ "name" .= byteStringToJSON (streamSourceInfoName source)
      , "filter_subject" .= byteStringToJSON (streamSourceInfoFilterSubject source)
      , "lag" .= streamSourceInfoLag source
      , "active" .= durationToNanoseconds (streamSourceInfoActive source)
      ]

instance FromJSON StreamSourceInfo where
  parseJSON =
    withObject "StreamSourceInfo" $ \value ->
      StreamSourceInfo
        <$> parseByteStringField value "name"
        <*> parseOptionalByteStringField value "filter_subject" .!= ""
        <*> value .: "lag"
        <*> parseDurationField value "active"

instance ToJSON DeleteStreamResponse where
  toJSON response =
    object
      [ "success" .= deleteStreamSuccess response
      ]

instance FromJSON DeleteStreamResponse where
  parseJSON =
    withObject "DeleteStreamResponse" $ \value ->
      DeleteStreamResponse
        <$> value .:? "success" .!= False

instance ToJSON PurgeStreamResponse where
  toJSON response =
    object
      [ "success" .= purgeStreamSuccess response
      , "purged" .= purgeStreamPurged response
      ]

instance FromJSON PurgeStreamResponse where
  parseJSON =
    withObject "PurgeStreamResponse" $ \value ->
      PurgeStreamResponse
        <$> value .:? "success" .!= False
        <*> value .:? "purged" .!= 0

instance ToJSON StreamListResponse where
  toJSON response =
    object
      [ "total" .= streamListTotal response
      , "offset" .= streamListOffset response
      , "limit" .= streamListLimit response
      , "streams" .= streamListStreams response
      ]

instance FromJSON StreamListResponse where
  parseJSON =
    withObject "StreamListResponse" $ \value ->
      StreamListResponse
        <$> value .: "total"
        <*> value .: "offset"
        <*> value .: "limit"
        <*> value .:? "streams" .!= []

instance ToJSON StreamNamesResponse where
  toJSON response =
    object
      [ "total" .= streamNamesTotal response
      , "offset" .= streamNamesOffset response
      , "limit" .= streamNamesLimit response
      , "streams" .= map byteStringToJSON (streamNamesStreams response)
      ]

instance FromJSON StreamNamesResponse where
  parseJSON =
    withObject "StreamNamesResponse" $ \value ->
      StreamNamesResponse
        <$> value .: "total"
        <*> value .: "offset"
        <*> value .: "limit"
        <*> parseOptionalByteStringListField value "streams" .!= []

instance FromJSON StreamMessage where
  parseJSON =
    withObject "StreamMessageResponse" $ \value ->
      value .: "message" >>= parseStoredMessage

instance ToJSON StreamMessage where
  toJSON message =
    object . catMaybes $
      [ Just (byteStringPair "subject" (streamMessageSubject message))
      , Just ("seq" .= streamMessageSequence message)
      , maybeBase64Pair "hdrs" (streamMessageHeadersRaw message)
      , maybeBase64Pair "data" (streamMessagePayload message)
      , Just ("time" .= streamMessageTime message)
      ]

instance FromJSON DeleteStreamMessageResponse where
  parseJSON =
    withObject "DeleteStreamMessageResponse" $ \value ->
      DeleteStreamMessageResponse
        <$> value .:? "success" .!= False

instance ToJSON DeleteStreamMessageResponse where
  toJSON response =
    object
      [ "success" .= deleteStreamMessageSuccess response
      ]

parseStoredMessage :: Value -> Parser StreamMessage
parseStoredMessage =
  withObject "StreamMessage" $ \value ->
    StreamMessage
      <$> parseByteStringField value "subject"
      <*> value .: "seq"
      <*> parseOptionalBase64Field value "hdrs"
      <*> parseOptionalBase64Field value "data"
      <*> value .: "time"

durationToNanoseconds :: NominalDiffTime -> Integer
durationToNanoseconds =
  diffTimeToNanoseconds

nanosecondsToDuration :: Integer -> NominalDiffTime
nanosecondsToDuration nanoseconds =
  fromRational (toRational nanoseconds / 1000000000)

byteStringPair :: Key -> BS.ByteString -> Pair
byteStringPair key value = key .= byteStringToJSON value

byteStringListPair :: Key -> [BS.ByteString] -> Pair
byteStringListPair key = (key .=) . map byteStringToJSON

maybeByteStringPair :: Key -> Maybe BS.ByteString -> Maybe Pair
maybeByteStringPair key = fmap (byteStringPair key)

maybeBase64Pair :: Key -> Maybe BS.ByteString -> Maybe Pair
maybeBase64Pair key =
  fmap ((key .=) . byteStringToJSON . Base64.encode)

maybeByteStringListPair :: Key -> Maybe [BS.ByteString] -> Maybe Pair
maybeByteStringListPair key = fmap (byteStringListPair key)

maybePair :: ToJSON value => Key -> Maybe value -> Maybe Pair
maybePair key = fmap (key .=)

maybeDurationPair :: Key -> Maybe NominalDiffTime -> Maybe Pair
maybeDurationPair key = fmap ((key .=) . durationToNanoseconds)

parseByteStringField :: Object -> Key -> Parser BS.ByteString
parseByteStringField value key =
  value .: key >>= parseByteString

parseOptionalByteStringField :: Object -> Key -> Parser (Maybe BS.ByteString)
parseOptionalByteStringField value key = do
  byteString <- value .:? key
  traverse parseByteString byteString

parseOptionalBase64Field :: Object -> Key -> Parser (Maybe BS.ByteString)
parseOptionalBase64Field value key = do
  encoded <- value .:? key
  traverse parseBase64 encoded
  where
    parseBase64 jsonValue = do
      bytes <- parseByteString jsonValue
      case Base64.decode bytes of
        Left err      -> fail err
        Right decoded -> pure decoded

parseOptionalByteStringListField :: Object -> Key -> Parser (Maybe [BS.ByteString])
parseOptionalByteStringListField value key = do
  byteStrings <- value .:? key
  traverse (traverse parseByteString) byteStrings

parseDurationField :: Object -> Key -> Parser NominalDiffTime
parseDurationField value key =
  nanosecondsToDuration <$> value .: key

parseOptionalDurationField :: Object -> Key -> Parser (Maybe NominalDiffTime)
parseOptionalDurationField value key =
  fmap nanosecondsToDuration <$> value .:? key