hslogstash 0.1.0.0 → 0.2.0
raw patch · 4 files changed
+161/−22 lines, 4 filesdep +attoparsecdep +conduitdep +hedisPVP ok
version bump matches the API change (PVP)
Dependencies added: attoparsec, conduit, hedis, http-conduit, text-format, time, transformers
API changes (from Hackage documentation)
- Logstash.Message: logStashContent :: LogstashMessage -> Text
- Logstash.Message: logStashFields :: LogstashMessage -> Value
- Logstash.Message: logStashSource :: LogstashMessage -> Text
- Logstash.Message: logStashTags :: LogstashMessage -> [Text]
- Logstash.Message: logStashType :: LogstashMessage -> Text
+ Data.Conduit.ElasticSearch: esSink :: MonadResource m => Maybe (Request m) -> ByteString -> Int -> Conduit LogstashMessage m (Either (LogstashMessage, Value) Value)
+ Data.Conduit.Redis: redisSource :: MonadResource m => HostName -> Int -> ByteString -> Source m ByteString
+ Logstash.Message: instance Eq LogstashMessage
+ Logstash.Message: logstashContent :: LogstashMessage -> Text
+ Logstash.Message: logstashFields :: LogstashMessage -> Value
+ Logstash.Message: logstashSource :: LogstashMessage -> Text
+ Logstash.Message: logstashTags :: LogstashMessage -> [Text]
+ Logstash.Message: logstashTime :: LogstashMessage -> Maybe UTCTime
+ Logstash.Message: logstashTimestamp :: UTCTime -> Text
+ Logstash.Message: logstashType :: LogstashMessage -> Text
+ Logstash.Message: parseLogstashTime :: Text -> Maybe UTCTime
- Logstash.Message: LogstashMessage :: Text -> Text -> [Text] -> Value -> Text -> LogstashMessage
+ Logstash.Message: LogstashMessage :: Text -> Text -> [Text] -> Value -> Text -> Maybe UTCTime -> LogstashMessage
Files
- Data/Conduit/ElasticSearch.hs +57/−0
- Data/Conduit/Redis.hs +30/−0
- Logstash/Message.hs +62/−15
- hslogstash.cabal +12/−7
+ Data/Conduit/ElasticSearch.hs view
@@ -0,0 +1,57 @@+{-| This module exports "Conduit" interfaces to ElasticSearch. It is+totally experimental.+-}+module Data.Conduit.ElasticSearch (esSink) where++import Prelude hiding (catch)+import Control.Exception+import Data.Conduit+import qualified Data.Conduit.List as CL+import Network.HTTP.Conduit+import Data.Aeson+import qualified Data.ByteString as BS+import qualified Data.ByteString.Lazy as BSL+import Data.Time+import qualified Data.Text.Lazy.Encoding as E+import Data.Text.Format (format,left)+import Logstash.Message+import Control.Monad.IO.Class+import Control.Concurrent (threadDelay)+import qualified Data.HashMap.Strict as HM++safeQuery :: Request (ResourceT IO) -> IO (Response BSL.ByteString)+safeQuery req = catch (withManager $ httpLbs req) (\e -> print (e :: SomeException) >> threadDelay 500000 >> safeQuery req)++-- | Takes JSONifiable values, and returns the result of the ES request+-- along with the value in case of errors, or ES's values in case of+-- success+esSink :: (MonadResource m) => Maybe (Request m) -- ^ Defaults parameters for the http request to ElasticSearch. Use "Nothing" for defaults.+ -> BS.ByteString -- ^ Hostname of the ElasticSearch server+ -> Int -- ^ Port of the HTTP interface (usually 9200)+ -> Conduit LogstashMessage m (Either (LogstashMessage, Value) Value)+esSink r h p = CL.mapM doIndexA+ where+ defR1 = case r of+ Just x -> x+ Nothing -> def+ defR2 = defR1 { host = h+ , port = p+ , method = "POST"+ , checkStatus = (\_ _ -> Nothing)+ }+ doIndexA :: (MonadResource m) => LogstashMessage -> m (Either (LogstashMessage, Value) Value)+ doIndexA input =+ case logstashTime input of+ Nothing -> return $! Left (input, object [ "error" .= String "Time was not supplied" ])+ Just (UTCTime day _) -> do+ let (y,m,d) = toGregorian day+ req = defR2 { path = BSL.toStrict (E.encodeUtf8 (format "/logstash-{}-{}-{}/{}/" (y, left 2 '0' m, left 2 '0' d, logstashType input)))+ , requestBody = RequestBodyLBS (encode input)+ }+ res <- liftIO $ safeQuery req+ case decode (responseBody res) of+ Just (Object hh) -> case HM.lookup "ok" hh of+ Just (Bool True) -> return $! Right (Object hh)+ _ -> return $! Left (input, Object hh)+ Just j -> return $! Left (input, j)+ Nothing -> return $! Left (input, object [ "error" .= String "Could not decode", "content" .= responseBody res ])
+ Data/Conduit/Redis.hs view
@@ -0,0 +1,30 @@+{-| Quick conduit for reading from Redis lists. Not tested much, and probably quite slow.+-}+module Data.Conduit.Redis where++import Data.Conduit+import Data.Conduit.Util+import qualified Data.ByteString.Char8 as BS+import Network+import Database.Redis hiding (String, decode)+import Control.Monad (void)+import Control.Monad.IO.Class (liftIO)++redisSource :: (MonadResource m) => HostName -- ^ Hostname of the Redis server+ -> Int -- ^ Port of the Redis server (usually 6379)+ -> BS.ByteString -- ^ Name of the list+ -> Source m BS.ByteString+redisSource h p list =+ let cinfo = defaultConnectInfo { connectHost = h, connectPort = PortNumber $ fromIntegral p }+ pull = do+ o <- blpop [list] 0+ case o of+ Right (Just (_,k)) -> return k+ _ -> pull+ in sourceStateIO (connect cinfo)+ (\conn -> runRedis conn (void quit))+ (\conn -> do+ o <- liftIO $ runRedis conn pull+ return (StateOpen conn o)+ )+
Logstash/Message.hs view
@@ -2,21 +2,27 @@ import Data.Aeson import qualified Data.Text as T+import qualified Data.Text.Lazy as TL import Control.Applicative import Control.Monad import qualified Data.HashMap.Strict as HM import qualified Data.Vector as V+import Data.Time+import Data.Text.Format +import Data.Attoparsec.Text+ {-| The Logstash message, as described in <https://github.com/logstash/logstash/wiki/logstash's-internal-message-format>. Please not there is no timestamp, as the logstash server will add it. -} data LogstashMessage = LogstashMessage- { logStashType :: T.Text- , logStashSource :: T.Text- , logStashTags :: [T.Text]- , logStashFields :: Value- , logStashContent :: T.Text- } deriving (Show)+ { logstashType :: T.Text+ , logstashSource :: T.Text+ , logstashTags :: [T.Text]+ , logstashFields :: Value+ , logstashContent :: T.Text+ , logstashTime :: Maybe UTCTime+ } deriving (Show, Eq) instance FromJSON LogstashMessage where parseJSON (Object v) = LogstashMessage@@ -25,22 +31,58 @@ <*> v .: "@tags" <*> v .: "@fields" <*> v .: "@message"+ <*> v .: "@timestamp" parseJSON _ = mzero {-| As the name implies, this creates a dummy Logstash message, only updating the message field. -} emptyLSMessage :: T.Text -> LogstashMessage-emptyLSMessage m = LogstashMessage "empty" "dummy" [] Null m+emptyLSMessage m = LogstashMessage "empty" "dummy" [] (object []) m Nothing instance ToJSON LogstashMessage where- toJSON (LogstashMessage ty s ta f c) = object [ "@type" .= ty- , "@source" .= s- , "@tags" .= ta- , "@fields" .= f- , "@message" .= c- ]+ toJSON (LogstashMessage ty s ta f c ts) = object $ [ "@type" .= ty+ , "@source" .= s+ , "@tags" .= ta+ , "@fields" .= f+ , "@message" .= c+ ] ++ case ts of+ Nothing -> []+ Just t -> [ "@timestamp" .= t ] +-- | This formats an UTCTime in what logstash expects+logstashTimestamp :: UTCTime -> T.Text+logstashTimestamp (UTCTime d t) = TL.toStrict $! format "{}-{}-{}T{}:{}:{}.{}Z" (year, tc month, tc day, tc hours, tc minutes, tc seconds, left 3 '0' imicro)+ where+ tc = left 2 '0'+ reduce :: Int -> Int -> (Int, Int)+ reduce a b = (a `mod` b, a `div` b)+ (year, month, day) = toGregorian d+ (fseconds, micro) = properFraction t+ imicro = truncate (micro * 1000) :: Int+ (seconds, fminutes) = reduce fseconds 60+ (minutes, hours) = reduce fminutes 60++-- | This parses the logstash format+parseLogstashTime :: T.Text -> Maybe UTCTime+parseLogstashTime t = case parseOnly prs t of+ Right r -> Just r+ Left _ -> Nothing+ where+ prs = do+ ye <- decimal <* char '-' :: Parser Integer+ mo <- decimal <* char '-' :: Parser Int+ da <- decimal <* char 'T' :: Parser Int+ ho <- decimal <* char ':' :: Parser Int+ mi <- decimal <* char ':' :: Parser Int+ se <- decimal <* char '.' :: Parser Int+ ms <- decimal <* char 'Z' :: Parser Int+ endOfInput+ let !seconds = ho*3600 + mi*60 + se+ !micro = fromIntegral ms / 1000+ !secs = secondsToDiffTime (fromIntegral seconds) + micro+ return $! UTCTime (fromGregorian ye mo da) secs+ {-| This will try to convert an arbitrary JSON value into a "LogstashMessage". -}@@ -53,18 +95,23 @@ Nothing -> Null mtags = case HM.lookup "@tags" m of Just (Array v) -> toTags (V.toList v)- Nothing -> Nothing+ _ -> Nothing mmsg = case HM.lookup "@message" m of Just (String x) -> x _ -> ""+ mts = case HM.lookup "@timestamp" m of+ Just (String u) -> parseLogstashTime u+ _ -> Nothing toTags :: [Value] -> Maybe [T.Text] toTags v = let isString (String _) = True isString _ = False toText (String x) = x+ toText _ = "" in if null (filter (not . isString) v) then Just (map toText v) else Nothing in case (mtype, msrc, mtags) of- (Just (String t), Just (String s), Just tags) -> Just $ LogstashMessage t s tags mflds mmsg+ (Just (String t), Just (String s), Just tags) -> Just $ LogstashMessage t s tags mflds mmsg mts _ -> Nothing+value2logstash _ = Nothing
hslogstash.cabal view
@@ -2,20 +2,25 @@ -- documentation, see http://haskell.org/cabal/users-guide/ name: hslogstash-version: 0.1.0.0-synopsis: A library to write structured messages to a logstash server.--- description: +version: 0.2.0+synopsis: A library to work with, or as, a logstash server+description: This library contains a few modules that let you work with Logstash messages, read them from a Redis list, store them into Elasticsearch, and more. license: BSD3 license-file: LICENSE author: Simon Marechal maintainer: bartavelle@gmail.com -- copyright: -category: System+category: System, Logstash build-type: Simple cabal-version: >=1.8 +source-repository head+ type: git+ location: git://github.com/bartavelle/hslogstash.git+ library- exposed-modules: Logstash.Message, Logstash.IO- extensions: OverloadedStrings+ exposed-modules: Logstash.Message, Logstash.IO, Data.Conduit.Redis, Data.Conduit.ElasticSearch+ extensions: OverloadedStrings, BangPatterns+ ghc-options: -Wall -- other-modules: - build-depends: base <5, aeson, network, bytestring, text, vector, unordered-containers+ build-depends: base <5, aeson, network, bytestring, text, vector, unordered-containers, time, text-format, attoparsec, hedis, conduit, transformers, http-conduit