amqp-utils 0.3.0.0 → 0.3.2.0
raw patch · 8 files changed
+1065/−510 lines, 8 filesdep +hinotifydep +magicdep +unixnew-component:exe:agitprop
Dependencies added: hinotify, magic, unix
Files
- ChangeLog.md +15/−0
- Network/AMQP/Utils/Connection.hs +72/−0
- Network/AMQP/Utils/Helpers.hs +185/−0
- Network/AMQP/Utils/Options.hs +366/−0
- README.md +117/−1
- agitprop.hs +149/−0
- amqp-utils.cabal +34/−2
- konsum.hs +127/−507
ChangeLog.md view
@@ -1,5 +1,20 @@ # Revision history for haskell-amqp-utils +## 0.3.2.0 -- 2018-07-04++* agitprop, a publisher+* optional publisher confirms+* hotfolder mode, file magic+* several options++## 0.3.0.2 -- 2018-04-24++* use ciphersuite\_default++## 0.3.0.1 -- 2018-03-04++* don't let the thread sleep too long+ ## 0.3.0.0 -- 2017-11-21 * add nix with `amqp_0_18_1`
+ Network/AMQP/Utils/Connection.hs view
@@ -0,0 +1,72 @@+{-# LANGUAGE OverloadedStrings #-}++module Network.AMQP.Utils.Connection where++import qualified Data.ByteString as B+import Data.Default.Class+import qualified Data.Text as T+import Network.AMQP+import Network.AMQP.Utils.Helpers+import Network.AMQP.Utils.Options+import qualified Network.Connection as N+import Network.TLS+import Network.TLS.Extra+import System.X509++-- | opens a connection and a channel+connect :: Args -> IO (Connection, Channel)+connect args = do+ printparam' "server" $ server args+ printparam' "port" $ show $ port args+ printparam' "vhost" $ vHost args+ printparam "connection_name" $ connectionName args+ globalCertificateStore <- getSystemCertificateStore+ let myTLS =+ N.TLSSettings+ (defaultParamsClient "" B.empty)+ { clientShared =+ def+ { sharedValidationCache = def+ , sharedCAStore = globalCertificateStore+ }+ , clientSupported = def {supportedCiphers = ciphersuite_default}+ , clientHooks =+ def {onCertificateRequest = myCert (cert args) (key args)}+ }+ conn <-+ openConnection''+ defaultConnectionOpts+ { coAuth =+ [ SASLMechanism "EXTERNAL" B.empty Nothing+ , plain (T.pack (user args)) (T.pack (pass args))+ ]+ , coVHost = T.pack $ vHost args+ , coTLSSettings =+ if (tls args)+ then Just (TLSCustom myTLS)+ else Nothing+ , coServers = [(server args, fromIntegral $ port args)]+ , coHeartbeatDelay = fmap fromIntegral $ heartBeat args+ , coName = fmap T.pack $ connectionName args+ }+ chan <- openChannel conn+ return (conn, chan)++-- addChannelExceptionHandler chan+-- (\exception -> closeConnection conn >>+-- printparam' "exiting" (show exception) >>+-- killThread tid)+--+-- -- noop sharedValidationCache, handy when debugging+-- noValidation :: ValidationCache+-- noValidation = ValidationCache+-- (\_ _ _ -> return ValidationCachePass)+-- (\_ _ _ -> return ())+-- | provides the TLS client certificate+myCert :: Maybe FilePath -> Maybe FilePath -> t -> IO (Maybe Credential)+myCert (Just cert') (Just key') _ = do+ result <- credentialLoadX509 cert' key'+ case result of+ Left x -> printparam' "ERROR" x >> return Nothing+ Right x -> return $ Just x+myCert _ _ _ = return Nothing
+ Network/AMQP/Utils/Helpers.hs view
@@ -0,0 +1,185 @@+module Network.AMQP.Utils.Helpers where++import qualified Data.ByteString.Lazy.Char8 as BL+import qualified Data.Map as M+import qualified Data.Text as T+import Data.Time+import Data.Time.Clock.POSIX+import Network.AMQP+import Network.AMQP.Types+import System.IO+import Data.Maybe+import Data.List+import Control.Monad++-- | log cmdline options+listToMaybeUnwords :: [String] -> Maybe String+listToMaybeUnwords [] = Nothing+listToMaybeUnwords x = Just $ unwords x++-- | Strings or ByteStrings with label, oder nothing at all+printwithlabel :: String -> Maybe (IO ()) -> IO ()+printwithlabel _ Nothing = return ()+printwithlabel labl (Just i) = do+ mapM_ putStr [" --- ", labl, ": "]+ i+ hFlush stdout++-- | optional parameters+printparam :: String -> Maybe String -> IO ()+printparam labl ms = printwithlabel labl $ fmap putStrLn ms++-- | required parameters+printparam' :: String -> String -> IO ()+printparam' d s = printparam d (Just s)++-- | head chars of body+printbody :: (String, Maybe BL.ByteString) -> IO ()+printbody (labl, ms) =+ printwithlabel labl $ fmap (\s -> putStrLn "" >> BL.putStrLn s) ms++-- | log marker+hr :: String -> IO ()+hr x = putStrLn hr' >> hFlush stdout+ where+ hr' = take 72 $ (take 25 hr'') ++ " " ++ x ++ " " ++ hr''+ hr'' = repeat '-'++formatheaders :: ((T.Text, FieldValue) -> [a]) -> FieldTable -> [a]+formatheaders f (FieldTable ll) = concat $ map f $ M.toList ll++-- | log formatting+fieldshow :: (T.Text, FieldValue) -> String+fieldshow (k, v) = "\n " ++ T.unpack k ++ ": " ++ valueshow v++-- | callback cmdline formatting+fieldshow' :: (T.Text, FieldValue) -> [String]+fieldshow' (k, v) = ["-h", T.unpack k ++ "=" ++ valueshow v]++-- | showing a FieldValue+valueshow :: FieldValue -> String+valueshow (FVString value) = T.unpack value+valueshow (FVInt32 value) = show value+valueshow value = show value++-- | skip showing body head if binary type+isimage :: Maybe String -> Bool+isimage Nothing = False+isimage (Just ctype)+ | isPrefixOf "application/xml" ctype = False+ | isPrefixOf "application/json" ctype = False+ | otherwise = any (flip isPrefixOf ctype) ["application", "image"]++-- | show the first bytes of message body+anriss' :: Maybe Int -> BL.ByteString -> BL.ByteString+anriss' x =+ case x of+ Nothing -> id+ Just y -> BL.take (fromIntegral y)++-- | callback cmdline with optional parameters+printopt :: (String, Maybe String) -> [String]+printopt (_, Nothing) = []+printopt (opt, Just s) = [opt, s]++-- | prints header and head on STDOUT and returns cmdline options to callback+printmsg :: (Message, Envelope) -> Maybe Int -> ZonedTime -> IO [String]+printmsg (msg, envi) anR now = do+ mapM_+ (uncurry printparam)+ [ ("routing key", rkey)+ , ("message-id", messageid)+ , ("headers", headers)+ , ("content-type", contenttype)+ , ("content-encoding", contentencoding)+ , ("redelivered", redeliv)+ , ("timestamp", timestamp'')+ , ("time now", now')+ , ("size", size)+ , ("priority", prio)+ , ("type", mtype)+ , ("user id", muserid)+ , ("application id", mappid)+ , ("cluster id", mclusterid)+ , ("reply to", mreplyto)+ , ("correlation id", mcorrid)+ , ("expiration", mexp)+ , ("delivery mode", mdelivmode)+ ]+ printbody (label, anriss)+ return $+ concat+ (map+ printopt+ [ ("-r", rkey)+ , ("-m", contenttype)+ , ("-e", contentencoding)+ , ("-i", messageid)+ , ("-t", timestamp)+ , ("-p", prio)+ ] +++ maybeToList headers')+ where+ headers = fmap (formatheaders fieldshow) $ msgHeaders msg+ headers' = fmap (formatheaders fieldshow') $ msgHeaders msg+ body = msgBody msg+ anriss =+ if isimage contenttype+ then Nothing+ else Just (anriss' anR body) :: Maybe BL.ByteString+ anriss'' = maybe "" (\a -> "first " ++ (show a) ++ " bytes of ") anR+ label = anriss'' ++ "body"+ contenttype = fmap T.unpack $ msgContentType msg+ contentencoding = fmap T.unpack $ msgContentEncoding msg+ rkey = Just . T.unpack $ envRoutingKey envi+ messageid = fmap T.unpack $ msgID msg+ prio = fmap show $ msgPriority msg+ mtype = fmap show $ msgType msg+ muserid = fmap show $ msgUserID msg+ mappid = fmap show $ msgApplicationID msg+ mclusterid = fmap show $ msgClusterID msg+ mreplyto = fmap show $ msgReplyTo msg+ mcorrid = fmap show $ msgCorrelationID msg+ mexp = fmap show $ msgExpiration msg+ mdelivmode = fmap show $ msgDeliveryMode msg+ size = Just . show $ BL.length body+ redeliv =+ if envRedelivered envi+ then Just "YES"+ else Nothing+ tz = zonedTimeZone now+ nowutc = zonedTimeToUTCFLoor now+ msgtime = msgTimestamp msg+ msgtimeutc = fmap (posixSecondsToUTCTime . realToFrac) msgtime+ timestamp = fmap show msgtime+ timediff = fmap (difftime nowutc) msgtimeutc+ now' =+ case timediff of+ Just "now" -> Nothing+ _ -> showtime tz $ Just nowutc+ timestamp' = showtime tz msgtimeutc+ timestamp'' =+ liftM3+ (\a b c -> a ++ " (" ++ b ++ ") (" ++ c ++ ")")+ timestamp+ timestamp'+ timediff++-- | timestamp conversion+zonedTimeToUTCFLoor :: ZonedTime -> UTCTime+zonedTimeToUTCFLoor x =+ posixSecondsToUTCTime $+ realToFrac ((floor . utcTimeToPOSIXSeconds . zonedTimeToUTC) x :: Timestamp)++-- | show the timestamp+showtime :: TimeZone -> Maybe UTCTime -> Maybe String+showtime tz = fmap (show . (utcToZonedTime tz))++-- | show difference between two timestamps+difftime :: UTCTime -> UTCTime -> String+difftime now msg+ | now == msg = "now"+ | now > msg = diff ++ " ago"+ | otherwise = diff ++ " in the future"+ where+ diff = show (diffUTCTime now msg)
+ Network/AMQP/Utils/Options.hs view
@@ -0,0 +1,366 @@+module Network.AMQP.Utils.Options where++import Data.Default.Class+import qualified Data.Map as M+import Data.Maybe+import Data.Text (Text, pack)+import Data.Version (showVersion)+import Network.AMQP+import Network.AMQP.Types+import Paths_amqp_utils (version)+import System.Console.GetOpt++-- | A data type for our options+data Args = Args+ { server :: String+ , port :: Int+ , tls :: Bool+ , vHost :: String+ , currentExchange :: String+ , bindings :: [(String, String)]+ , rKey :: String+ , anRiss :: Maybe Int+ , fileProcess :: Maybe String+ , qName :: Maybe String+ , cert :: Maybe String+ , key :: Maybe String+ , user :: String+ , pass :: String+ , preFetch :: Int+ , heartBeat :: Maybe Int+ , tempDir :: Maybe String+ , additionalArgs :: [String]+ , connectionName :: Maybe String+ , tmpQName :: String+ , inputFile :: String+ , lineMode :: Bool+ , confirm :: Bool+ , msgid :: Maybe Text+ , msgtype :: Maybe Text+ , userid :: Maybe Text+ , appid :: Maybe Text+ , clusterid :: Maybe Text+ , contenttype :: Maybe Text+ , contentencoding :: Maybe Text+ , replyto :: Maybe Text+ , prio :: Maybe Octet+ , corrid :: Maybe Text+ , msgexp :: Maybe Text+ , msgheader :: Maybe FieldTable+ , fnheader :: [String]+ , suffix :: [String]+ , magic :: Bool+ , persistent :: Maybe DeliveryMode+ }++instance Default Args where+ def =+ Args+ "localhost"+ 5672+ False+ "/"+ "default"+ []+ ""+ Nothing+ Nothing+ Nothing+ Nothing+ Nothing+ "guest"+ "guest"+ 1+ Nothing+ Nothing+ []+ Nothing+ ""+ "/dev/stdin"+ False+ False+ Nothing+ Nothing+ Nothing+ Nothing+ Nothing+ Nothing+ Nothing+ Nothing+ Nothing+ Nothing+ Nothing+ Nothing+ []+ []+ False+ Nothing++-- | Common options+cOptions :: [OptDescr (Args -> Args)]+cOptions =+ [ Option+ ['o']+ ["server"]+ (ReqArg (\s o -> o {server = s}) "SERVER")+ ("AMQP Server (default: " ++ server def ++ ")")+ , Option+ ['y']+ ["vhost"]+ (ReqArg (\s o -> o {vHost = s}) "VHOST")+ ("AMQP Virtual Host (default: " ++ vHost def ++ ")")+ , Option+ ['x']+ ["exchange"]+ (ReqArg (\s o -> o {currentExchange = s}) "EXCHANGE")+ ("AMQP Exchange (default: default)")+ , Option+ ['Q']+ ["qname"]+ (ReqArg (\s o -> o {tmpQName = s}) "TEMPQNAME")+ "Name for temporary exclusive Queue"+ , Option+ ['p']+ ["port"]+ (ReqArg (\s o -> o {port = read s}) "PORT")+ ("Server Port Number (default: " ++ show (port def) ++ ")")+ , Option+ ['T']+ ["tls"]+ (NoArg (\o -> o {tls = not (tls o)}))+ ("Toggle TLS (default: " ++ show (tls def) ++ ")")+ , Option+ ['q']+ ["queue"]+ (ReqArg (\s o -> o {qName = Just s}) "QUEUENAME")+ "Ignore Exchange and bind to existing Queue"+ , Option+ ['c']+ ["cert"]+ (ReqArg (\s o -> o {cert = Just s}) "CERTFILE")+ ("TLS Client Certificate File")+ , Option+ ['k']+ ["key"]+ (ReqArg (\s o -> o {key = Just s}) "KEYFILE")+ ("TLS Client Private Key File")+ , Option+ ['U']+ ["user"]+ (ReqArg (\s o -> o {user = s}) "USERNAME")+ ("Username for Auth")+ , Option+ ['P']+ ["pass"]+ (ReqArg (\s o -> o {pass = s}) "PASSWORD")+ ("Password for Auth")+ , Option+ ['s']+ ["heartbeats"]+ (ReqArg (\s o -> o {heartBeat = (Just (read s))}) "INT")+ "heartbeat interval (0=disable, default: set by server)"+ , Option+ ['n']+ ["name"]+ (ReqArg (\s o -> o {connectionName = Just s}) "NAME")+ "connection name, will be shown in RabbitMQ web interface"+ ]++-- | Options for konsum+kOptions :: [OptDescr (Args -> Args)]+kOptions =+ [ Option+ ['r']+ ["bindingkey"]+ (ReqArg+ (\s o -> o {bindings = (currentExchange o, s) : (bindings o)})+ "BINDINGKEY")+ ("AMQP binding key (default: #)")+ , Option+ ['X']+ ["execute"]+ (OptArg+ (\s o ->+ o+ { fileProcess = Just (fromMaybe callback s)+ , tempDir = Just (fromMaybe "/tmp" (tempDir o))+ })+ "EXE")+ ("Callback Script File (implies -t) (-X without arg: " ++ callback ++ ")")+ , Option+ ['a']+ ["args"]+ (ReqArg (\s o -> o {additionalArgs = s : (additionalArgs o)}) "ARG")+ "additional argument for -X callback"+ , Option+ ['l']+ ["charlimit"]+ (ReqArg (\s o -> o {anRiss = Just (read s :: Int)}) "INT")+ "limit number of shown body chars (default: unlimited)"+ , Option+ ['t']+ ["tempdir", "target"]+ (OptArg (\s o -> o {tempDir = Just (fromMaybe "/tmp" s)}) "DIR")+ "tempdir (default: no file creation, -t without arg: /tmp)"+ , Option+ ['f']+ ["prefetch"]+ (ReqArg (\s o -> o {preFetch = read s}) "INT")+ ("Prefetch count. (0=unlimited, 1=off, default: " +++ show (preFetch def) ++ ")")+ ]++-- | Options for agitprop+aOptions :: [OptDescr (Args -> Args)]+aOptions =+ [ Option+ ['r']+ ["routingkey"]+ (ReqArg (\s o -> o {rKey = s}) "ROUTINGKEY")+ "AMQP routing key"+ , Option+ ['f']+ ["inputfile"]+ (ReqArg (\s o -> o {inputFile = s}) "INPUTFILE")+ ("Message input file (default: " ++ (inputFile def) ++ ")")+ , Option+ ['l']+ ["linemode"]+ (NoArg (\o -> o {lineMode = not (lineMode o)}))+ ("Toggle line-by-line mode (default: " ++ show (lineMode def) ++ ")")+ , Option+ ['C']+ ["confirm"]+ (NoArg (\o -> o {confirm = not (confirm o)}))+ ("Toggle confirms (default: " ++ show (confirm def) ++ ")")+ , Option+ []+ ["msgid"]+ (ReqArg (\s o -> o {msgid = Just $ pack s}) "ID")+ "Message ID"+ , Option+ []+ ["type"]+ (ReqArg (\s o -> o {msgtype = Just $ pack s}) "TYPE")+ "Message Type"+ , Option+ []+ ["userid"]+ (ReqArg (\s o -> o {userid = Just $ pack s}) "USERID")+ "Message User-ID"+ , Option+ []+ ["appid"]+ (ReqArg (\s o -> o {appid = Just $ pack s}) "APPID")+ "Message App-ID"+ , Option+ []+ ["clusterid"]+ (ReqArg (\s o -> o {clusterid = Just $ pack s}) "CLUSTERID")+ "Message Cluster-ID"+ , Option+ []+ ["contenttype"]+ (ReqArg (\s o -> o {contenttype = Just $ pack s}) "CONTENTTYPE")+ "Message Content-Type"+ , Option+ []+ ["contentencoding"]+ (ReqArg (\s o -> o {contentencoding = Just $ pack s}) "CONTENTENCODING")+ "Message Content-Encoding"+ , Option+ []+ ["replyto"]+ (ReqArg (\s o -> o {replyto = Just $ pack s}) "REPLYTO")+ "Message Reply-To"+ , Option+ []+ ["prio"]+ (ReqArg (\s o -> o {prio = Just $ read s}) "PRIO")+ "Message Priority"+ , Option+ []+ ["corrid"]+ (ReqArg (\s o -> o {corrid = Just $ pack s}) "CORRID")+ "Message CorrelationID"+ , Option+ []+ ["exp"]+ (ReqArg (\s o -> o {msgexp = Just $ pack s}) "EXP")+ "Message Expiration"+ , Option+ ['h']+ ["header"]+ (ReqArg (\s o -> o {msgheader = addheader (msgheader o) s}) "HEADER=VALUE")+ "Message Headers"+ , Option+ ['F']+ ["fnheader"]+ (ReqArg (\s o -> o {fnheader = s : (fnheader o)}) "HEADERNAME")+ "Put filename into this header"+ , Option+ ['S']+ ["suffix"]+ (ReqArg (\s o -> o {suffix = s : (suffix o)}) "SUFFIX")+ "Allowed file suffixes in hotfolder mode"+ , Option+ ['m']+ ["magic"]+ (NoArg (\o -> o {magic = not (magic o)}))+ ("Toggle setting content-type and -encoding from file contents (default: " +++ show (magic def) ++ ")")+ , Option+ ['e']+ ["persistent"]+ (NoArg (\o -> o {persistent = Just Persistent}))+ "Set persistent delivery"+ , Option+ ['E']+ ["nonpersistent"]+ (NoArg (\o -> o {persistent = Just NonPersistent}))+ "Set nonpersistent delivery"+ ]++-- | Options for the executables+options :: String -> [OptDescr (Args -> Args)]+options "konsum" = kOptions ++ cOptions+options "agitprop" = aOptions ++ cOptions+options _ = cOptions++-- | Add a header with a String value+addheader :: Maybe FieldTable -> String -> Maybe FieldTable+addheader Nothing string = Just $ FieldTable $ M.singleton (k string) (v string)+addheader (Just (FieldTable oldheader)) string =+ Just $ FieldTable $ M.insert (k string) (v string) oldheader++k :: String -> Text+k s = pack $ takeWhile (/= '=') s++v :: String -> FieldValue+v s = FVString $ pack $ tail $ dropWhile (/= '=') s++-- | 'parseargs' exename argstring+-- applies options onto argstring+parseargs :: String -> [String] -> IO Args+parseargs exename argstring =+ case getOpt Permute opt argstring of+ (o, [], []) -> return $ foldl (flip id) def o+ (_, _, errs) ->+ ioError $ userError $ concat errs ++ usageInfo (usage exename) opt+ where+ opt = options exename++-- | the default callback for the -X option+callback :: String+callback = "/usr/lib/haskell-amqp-utils/callback"++usage :: String -> String+usage exename =+ "\n\+ \amqp-utils " +++ (showVersion version) +++ "\n\n\+ \Usage:\n" +++ exename +++ " [options]\n\n\+ \Options:"
README.md view
@@ -1,2 +1,118 @@ # haskell-amqp-utils-generic Haskell AMQP consumer for use with RabbitMQ+generic Haskell AMQP commandline tools for use with RabbitMQ++## Overview+The package contains 2 binaries for commandline use.++- konsum, a generic consumer+- agitprop, a generic publisher++## konsum+### usage++ konsum [options]++ Options:+ -r BINDINGKEY --bindingkey=BINDINGKEY AMQP binding key (default: #)+ -X[EXE] --execute[=EXE] Callback Script File (implies -t) (-X without arg: /usr/lib/haskell-amqp-utils/callback)+ -a ARG --args=ARG additional argument for -X callback+ -l INT --charlimit=INT limit number of shown body chars (default: unlimited)+ -t[DIR] --tempdir[=DIR], --target[=DIR] tempdir (default: no file creation, -t without arg: /tmp)+ -f INT --prefetch=INT Prefetch count. (0=unlimited, 1=off, default: 1)+ -o SERVER --server=SERVER AMQP Server (default: localhost)+ -y VHOST --vhost=VHOST AMQP Virtual Host (default: /)+ -x EXCHANGE --exchange=EXCHANGE AMQP Exchange (default: default)+ -Q TEMPQNAME --qname=TEMPQNAME Name for temporary exclusive Queue+ -p PORT --port=PORT Server Port Number (default: 5672)+ -T --tls Toggle TLS (default: False)+ -q QUEUENAME --queue=QUEUENAME Ignore Exchange and bind to existing Queue+ -c CERTFILE --cert=CERTFILE TLS Client Certificate File+ -k KEYFILE --key=KEYFILE TLS Client Private Key File+ -U USERNAME --user=USERNAME Username for Auth+ -P PASSWORD --pass=PASSWORD Password for Auth+ -s INT --heartbeats=INT heartbeat interval (0=disable, default: set by server)+ -n NAME --name=NAME connection name, will be shown in RabbitMQ web interface++## examples++connect to localhost with default credentials and attach to a new temp+queue on exchange "default":++ ./konsum++Connect to a host with TLS on a custom port, authenticating with SSL+client certificate. On every received message a callback is spawned.+The message will be ACKed when the callback exits successfully. First+500 bytes of the message body are printed to stdout. Header infos are+always printed to stdout:++ konsum -o amqp.example.com -p 5673 -T -k amqp-key.pem -c amqp-crt.pem -y vhost -x exchange -X./callback.sh -a -c -a callback.config.sh -f 2 -r routing.key.# -l 500++Authenticate with user and pass. Generate a file for every message:++ konsum -o amqp.example.com -U user -P pass -q queue -t++Provide a custom CA cert for proving the server's identity via+enviroment:++ $ env SYSTEM_CERTIFICATE_PATH=/etc/amqp/cacert.crt ./konsum -T -y vhost -x exchange++Stop with ^C++## agitprop+### usage++ agitprop [options]++ Options:+ -r ROUTINGKEY --routingkey=ROUTINGKEY AMQP routing key+ -f INPUTFILE --inputfile=INPUTFILE Message input file (default: /dev/stdin)+ -l --linemode Toggle line-by-line mode (default: False)+ -C --confirm Toggle confirms (default: False)+ --msgid=ID Message ID+ --type=TYPE Message Type+ --userid=USERID Message User-ID+ --appid=APPID Message App-ID+ --clusterid=CLUSTERID Message Cluster-ID+ --contenttype=CONTENTTYPE Message Content-Type+ --contentencoding=CONTENTENCODING Message Content-Encoding+ --replyto=REPLYTO Message Reply-To+ --prio=PRIO Message Priority+ --corrid=CORRID Message CorrelationID+ --exp=EXP Message Expiration+ -h HEADER=VALUE --header=HEADER=VALUE Message Headers+ -F HEADERNAME --fnheader=HEADERNAME Put filename into this header+ -S SUFFIX --suffix=SUFFIX Allowed file suffixes in hotfolder mode+ -m --magic Toggle setting content-type and -encoding from file contents (default: False)+ -e --persistent Set persistent delivery+ -E --nonpersistent Set nonpersistent delivery+ -o SERVER --server=SERVER AMQP Server (default: localhost)+ -y VHOST --vhost=VHOST AMQP Virtual Host (default: /)+ -x EXCHANGE --exchange=EXCHANGE AMQP Exchange (default: default)+ -Q TEMPQNAME --qname=TEMPQNAME Name for temporary exclusive Queue+ -p PORT --port=PORT Server Port Number (default: 5672)+ -T --tls Toggle TLS (default: False)+ -q QUEUENAME --queue=QUEUENAME Ignore Exchange and bind to existing Queue+ -c CERTFILE --cert=CERTFILE TLS Client Certificate File+ -k KEYFILE --key=KEYFILE TLS Client Private Key File+ -U USERNAME --user=USERNAME Username for Auth+ -P PASSWORD --pass=PASSWORD Password for Auth+ -s INT --heartbeats=INT heartbeat interval (0=disable, default: set by server)+ -n NAME --name=NAME connection name, will be shown in RabbitMQ web interface++If INPUTFILE is a file, the file is sent as a message and the program+exits. If INPUTFILE is a directory, the directory is watched via inotify+and every file, which is written and closed or moved in gets sent,+optionally only files which match one or several SUFFIXes. Optionally+the file name is written into a message header named HEADERNAME.+Optionally Content-Type and Content-Encoding headers are set via magic+retrieved from file contents.++Line-by-line mode sends one message per INPUTFILE line.++### examples++Send a message containing a file and put the filename into a fileName+message header:++ agitprop -x amq.topic -r test -F fileName -f agitprop.hs
+ agitprop.hs view
@@ -0,0 +1,149 @@+{-# LANGUAGE OverloadedStrings #-}++import Control.Concurrent (threadDelay)+import qualified Control.Exception as X+import Control.Monad (forever)+import qualified Data.ByteString.Lazy.Char8 as BL+import Data.List (isSuffixOf)+import Data.Maybe+import qualified Data.Text as T+import Data.Time+import Data.Time.Clock.POSIX+import Data.Version (showVersion)+import Data.Word (Word64)+import Magic+import Network.AMQP+import Network.AMQP.Types+import Network.AMQP.Utils.Connection+import Network.AMQP.Utils.Helpers+import Network.AMQP.Utils.Options+import Paths_amqp_utils (version)+import System.Environment+import System.INotify+import qualified System.Posix.Files as F++main :: IO ()+main = do+ hr "starting"+ args <- getArgs >>= parseargs "agitprop"+ printparam' "client version" $ "amqp-utils " ++ (showVersion version)+ printparam' "routing key" $ rKey args+ isDir <- F.getFileStatus (inputFile args) >>= return . F.isDirectory+ if isDir+ then printparam' "hotfolder" $ inputFile args+ else printparam' "input file" $+ (inputFile args) +++ if (lineMode args)+ then " (line-by-line)"+ else ""+ (conn, chan) <- connect args+ printparam' "confirm mode" $ show $ confirm args+ if (confirm args)+ then do+ confirmSelect chan False+ addConfirmationListener chan confirmCallback+ else return ()+ let publishOneMsg = publishOneMsg' chan args+ X.catch+ (if isDir+ then do+ inotify <- initINotify+ wd <-+ addWatch+ inotify+ [CloseWrite, MoveIn]+ (inputFile args)+ (handleEvent publishOneMsg (suffix args))+ hr $ "watching " ++ (inputFile args)+ _ <- forever $ threadDelay 1000000+ removeWatch wd+ else do+ hr $ "sending " ++ (inputFile args)+ messageFile <- BL.readFile (inputFile args)+ if (lineMode args)+ then mapM_ (publishOneMsg Nothing) (BL.lines messageFile)+ else publishOneMsg (Just (inputFile args)) messageFile)+ (\exception -> printparam' "exception" $ show (exception :: X.SomeException))+ -- all done. wait and close.+ if (confirm args)+ then waitForConfirms chan >>= return . show >> return ()+ else return ()+ closeConnection conn++-- | The handler for publisher confirms+confirmCallback :: (Word64, Bool, AckType) -> IO ()+confirmCallback (deliveryTag, isAll, ackType) =+ printparam'+ "confirmed"+ ((show deliveryTag) +++ (if isAll+ then " all "+ else " this ") +++ (show ackType))++-- | Hotfolder event handler+handleEvent ::+ (Maybe String -> BL.ByteString -> IO ()) -> [String] -> Event -> IO ()+-- just handle closewrite and movedin events+handleEvent f s (Closed False (Just x) True) = handleFile f s x+handleEvent f s (MovedIn False x _) = handleFile f s x+handleEvent _ _ _ = return ()++-- | Hotfolder file handler+handleFile ::+ (Maybe String -> BL.ByteString -> IO ()) -> [String] -> FilePath -> IO ()+handleFile _ _ ('.':_) = return () -- ignore hidden files+handleFile f s@(_:_) x =+ if any (flip isSuffixOf x) s+ then handleFile f [] x+ else return ()+handleFile f [] x =+ X.catch+ (hr ("sending " ++ x) >> BL.readFile x >>= f (Just x))+ (\exception ->+ printparam' "exception in handleFile" $+ show (exception :: X.SomeException))++-- | Publish one message with our settings+publishOneMsg' :: Channel -> Args -> Maybe String -> BL.ByteString -> IO ()+publishOneMsg' c a fn f = do+ (mtype, mencoding) <-+ if (magic a) && isJust fn+ then do+ m <- magicOpen [MagicMimeType]+ magicLoadDefault m+ t <- magicFile m (fromJust fn)+ magicSetFlags m [MagicMimeEncoding]+ e <- magicFile m (fromJust fn)+ return (Just (T.pack t), Just (T.pack e))+ else return ((contenttype a), (contentencoding a))+ now <- getCurrentTime >>= return . floor . utcTimeToPOSIXSeconds+ r <-+ publishMsg+ c+ (T.pack $ currentExchange a)+ (T.pack $ rKey a)+ newMsg+ { msgBody = f+ , msgDeliveryMode = persistent a+ , msgTimestamp = Just now+ , msgID = msgid a+ , msgType = msgtype a+ , msgUserID = userid a+ , msgApplicationID = appid a+ , msgClusterID = clusterid a+ , msgContentType = mtype+ , msgContentEncoding = mencoding+ , msgReplyTo = replyto a+ , msgPriority = prio a+ , msgCorrelationID = corrid a+ , msgExpiration = msgexp a+ , msgHeaders = substheader (fnheader a) fn $ msgheader a+ }+ printparam "sent" $ fmap show r+ where+ substheader ::+ [String] -> Maybe String -> Maybe FieldTable -> Maybe FieldTable+ substheader (s:r) (Just fname) old =+ substheader r (Just fname) (addheader old (s ++ "=" ++ fname))+ substheader _ _ old = old
amqp-utils.cabal view
@@ -1,6 +1,6 @@ name: amqp-utils -version: 0.3.0.0+version: 0.3.2.0 synopsis: Generic Haskell AMQP Consumer @@ -27,7 +27,7 @@ cabal-version: >=1.10 -Tested-With: GHC == 7.10.2, GHC == 8.0.2+Tested-With: GHC == 7.10.2, GHC == 8.0.2, GHC == 8.2.2 executable konsum main-is: konsum.hs@@ -46,7 +46,39 @@ ghc-options: -threaded -Wall default-language: Haskell98++ other-modules: Network.AMQP.Utils.Options,+ Network.AMQP.Utils.Helpers,+ Network.AMQP.Utils.Connection,+ Paths_amqp_utils++executable agitprop+ main-is: agitprop.hs+ build-depends: base >=4.6 && <5,+ containers,+ text,+ connection,+ data-default-class,+ time,+ process,+ bytestring,+ x509-system,+ tls,+ amqp >=0.17,+ unix >= 2.7,+ hinotify >= 0.3.8 && < 0.3.10,+ magic + ghc-options: -threaded -Wall+ + default-language: Haskell98++ other-modules: Network.AMQP.Utils.Options,+ Network.AMQP.Utils.Helpers,+ Network.AMQP.Utils.Connection,+ Paths_amqp_utils++ source-repository head type: git location: git://github.com/woffs/haskell-amqp-utils
konsum.hs view
@@ -1,3 +1,5 @@+{-# LANGUAGE OverloadedStrings #-}+ -- generic amqp consumer -- -- compile:@@ -13,531 +15,149 @@ -- -- Stop with ^C -import Paths_amqp_utils ( version )-import Data.Version ( showVersion )-import System.Environment-import System.Console.GetOpt-import System.IO-import System.Process-import System.Exit-import Control.Monad-import Control.Concurrent-import qualified Control.Exception as X-import Data.List-import Data.Maybe-import qualified Data.Text as T-import qualified Data.ByteString as B-import qualified Data.Map as M-import Network.AMQP-import Network.AMQP.Types-import Network.Connection-import Network.TLS-import Network.TLS.Extra-import System.X509+import Control.Concurrent+import qualified Control.Exception as X+import Control.Monad import qualified Data.ByteString.Lazy.Char8 as BL-import Data.Time-import Data.Time.Clock.POSIX-import Data.Default.Class+import qualified Data.Text as T+import Data.Time+import Data.Version (showVersion)+import Network.AMQP+import Network.AMQP.Utils.Connection+import Network.AMQP.Utils.Helpers+import Network.AMQP.Utils.Options+import Paths_amqp_utils (version)+import System.Environment+import System.Exit+import System.IO+import System.Process main :: IO () main = do- hr "starting"- tid <- myThreadId- args <- getArgs >>= parseargs- let addiArgs = reverse $ additionalArgs args- printparam' "client version" $ "amqp-utils " ++ (showVersion version)- printparam' "server" $ server args- printparam' "port" $ show $ port args- printparam' "vhost" $ vHost args- printparam "connection_name" $ connectionName args- globalCertificateStore <- getSystemCertificateStore- let myTLS = TLSSettings (defaultParamsClient "" B.empty)- { clientShared = def- { sharedValidationCache = def- , sharedCAStore = globalCertificateStore- }- , clientSupported = def- { supportedCiphers = ciphersuite_all- }- , clientHooks = def- { onCertificateRequest = myCert (cert args) (key args)- }- }- conn <- openConnection'' defaultConnectionOpts- { coAuth = [ SASLMechanism (T.pack "EXTERNAL") B.empty Nothing- , plain (T.pack (user args)) (T.pack (pass args))- ]- , coVHost = T.pack $ vHost args- , coTLSSettings = if (tls args ) then Just ( TLSCustom myTLS ) else Nothing- , coServers = [ (server args, fromIntegral $ port args) ]- , coHeartbeatDelay = fmap fromIntegral $ heartBeat args- , coName = fmap T.pack $ connectionName args- }- chan <- openChannel conn- addChannelExceptionHandler chan- (\exception -> closeConnection conn >>- printparam' "exiting" (show exception) >>- killThread tid)-- -- set prefetch- printparam' "prefetch" $- show $ preFetch args- qos chan 0 (fromIntegral $ preFetch args) False-- -- attach to given queue? or build exclusive queue and bind it?- queue <- maybe (tempQueue chan (tmpQName args) (bindings args) (currentExchange args))- (return)- (fmap T.pack (qName args))- printparam' "queue name" $ T.unpack queue-- printparam "shown body chars" $ fmap show $ anRiss args- printparam "temp dir" $ tempDir args- printparam "callback" $ fileProcess args- printparam "callback args" $ listToMaybeUnwords addiArgs-- -- subscribe to the queue- ctag <- consumeMsgs chan- queue- Ack- (myCallback (anRiss args)- (fileProcess args)- (tempDir args)- addiArgs- tid)- printparam' "consumer tag" $ T.unpack ctag- hr "entering main loop"-- X.catch (forever $ threadDelay 3600000000)- (\exception -> printparam' "exception" $- show (exception :: X.SomeException))- closeConnection conn---- -- noop sharedValidationCache, handy when debugging--- noValidation :: ValidationCache--- noValidation = ValidationCache--- (\_ _ _ -> return ValidationCachePass)--- (\_ _ _ -> return ())---- client certificate-myCert :: Maybe FilePath -> Maybe FilePath -> t -> IO (Maybe Credential)-myCert (Just cert') (Just key') _ = do- result <- credentialLoadX509 cert' key'- case result of- Left x -> printparam' "ERROR" x >> return Nothing- Right x -> return $ Just x-myCert _ _ _ = return Nothing+ hr "starting"+ tid <- myThreadId+ args <- getArgs >>= parseargs "konsum"+ let addiArgs = reverse $ additionalArgs args+ printparam' "client version" $ "amqp-utils " ++ (showVersion version)+ (conn, chan) <- connect args+ -- set prefetch+ printparam' "prefetch" $ show $ preFetch args+ qos chan 0 (fromIntegral $ preFetch args) False+ -- attach to given queue? or build exclusive queue and bind it?+ queue <-+ maybe+ (tempQueue chan (tmpQName args) (bindings args) (currentExchange args))+ (return)+ (fmap T.pack (qName args))+ printparam' "queue name" $ T.unpack queue+ printparam "shown body chars" $ fmap show $ anRiss args+ printparam "temp dir" $ tempDir args+ printparam "callback" $ fileProcess args+ printparam "callback args" $ listToMaybeUnwords addiArgs+ -- subscribe to the queue+ ctag <-+ consumeMsgs+ chan+ queue+ Ack+ (myCallback (anRiss args) (fileProcess args) (tempDir args) addiArgs tid)+ printparam' "consumer tag" $ T.unpack ctag+ hr "entering main loop"+ X.catch+ (forever $ threadDelay 5000000)+ (\exception -> printparam' "exception" $ show (exception :: X.SomeException))+ closeConnection conn --- exclusive temp queue+-- | exclusive temp queue tempQueue :: Channel -> String -> [(String, String)] -> String -> IO T.Text tempQueue chan tmpqname bindlist x = do- (q, _, _) <- declareQueue chan- newQueue { queueExclusive = True- , queueName = T.pack tmpqname- }- mapM_ (\(xchange, bkey) -> bindQueue chan q (T.pack xchange) (T.pack bkey) >>- printparam' "binding" (xchange ++ ":" ++ bkey))- (if null bindlist then [ (x, "#") ] else bindlist)- return q---- log cmdline options-listToMaybeUnwords :: [String] -> Maybe String-listToMaybeUnwords [] = Nothing-listToMaybeUnwords x = Just $ unwords x---- Strings or ByteStrings with label, oder nothing at all-printwithlabel :: String -> Maybe (IO ()) -> IO ()-printwithlabel _ Nothing =- return ()-printwithlabel labl (Just i) = do- mapM_ putStr [ " --- ", labl, ": " ]- i- hFlush stdout---- optional parameters-printparam :: String -> Maybe String -> IO ()-printparam labl ms = printwithlabel labl $- fmap putStrLn ms---- required parameters-printparam' :: String -> String -> IO ()-printparam' d s = printparam d (Just s)---- head chars of body-printbody :: (String, Maybe BL.ByteString) -> IO ()-printbody (labl, ms) = printwithlabel labl $- fmap (\s -> putStrLn "" >> BL.putStrLn s) ms---- options for everybody-data Args = Args { server :: String- , port :: Int- , tls :: Bool- , vHost :: String- , currentExchange :: String- , bindings :: [(String, String)]- , anRiss :: Maybe Int- , fileProcess :: Maybe String- , qName :: Maybe String- , cert :: Maybe String- , key :: Maybe String- , user :: String- , pass :: String- , preFetch :: Int- , heartBeat :: Maybe Int- , tempDir :: Maybe String- , additionalArgs :: [String]- , connectionName :: Maybe String- , tmpQName :: String- }--instance Default Args where- def = Args "localhost"- 5672- False- "/"- "default"- []- Nothing- Nothing- Nothing- Nothing- Nothing- "guest"- "guest"- 1- Nothing- Nothing- []- Nothing- ""--callback :: String-callback = "/usr/lib/haskell-amqp-utils/callback"--options :: [OptDescr (Args -> Args)]-options = [ Option [ 'o' ]- [ "server" ]- (ReqArg (\s o -> o { server = s }) "SERVER")- ("AMQP Server (default: " ++ server def ++ ")")- , Option [ 'y' ]- [ "vhost" ]- (ReqArg (\s o -> o { vHost = s }) "VHOST")- ("AMQP Virtual Host (default: " ++ vHost def ++ ")")- , Option [ 'x' ]- [ "exchange" ]- (ReqArg (\s o -> o { currentExchange = s }) "EXCHANGE")- ("AMQP Exchange (default: default)")- , Option [ 'r' ]- [ "bindingkey" ]- (ReqArg (\s o -> o { bindings = (currentExchange o, s) :- (bindings o)- })- "BINDINGKEY")- ("AMQP binding key (default: #)")- , Option [ 'Q' ]- [ "qname" ]- (ReqArg (\s o -> o { tmpQName = s }) "TEMPQNAME")- "Name for temporary exclusive Queue"- , Option [ 'X' ]- [ "execute" ]- (OptArg (\s o -> o { fileProcess = Just (fromMaybe callback s)- , tempDir = Just (fromMaybe "/tmp"- (tempDir o))- })- "EXE")- ("Callback Script File (implies -t) (-X without arg: " ++- callback ++ ")")- , Option [ 'p' ]- [ "port" ]- (ReqArg (\s o -> o { port = read s }) "PORT")- ("Server Port Number (default: " ++ show (port def) ++ ")")- , Option [ 'T' ]- [ "tls" ]- (NoArg (\o -> o { tls = not (tls o) }))- ("Toggle TLS (default: " ++ show (tls def) ++ ")")- , Option [ 'q' ]- [ "queue" ]- (ReqArg (\s o -> o { qName = Just s }) "QUEUENAME")- "Ignore Exchange and bind to existing Queue"- , Option [ 'c' ]- [ "cert" ]- (ReqArg (\s o -> o { cert = Just s }) "CERTFILE")- ("TLS Client Certificate File")- , Option [ 'k' ]- [ "key" ]- (ReqArg (\s o -> o { key = Just s }) "KEYFILE")- ("TLS Client Private Key File")- , Option [ 'U' ]- [ "user" ]- (ReqArg (\s o -> o { user = s }) "USERNAME")- ("Username for Auth")- , Option [ 'P' ]- [ "pass" ]- (ReqArg (\s o -> o { pass = s }) "PASSWORD")- ("Password for Auth")- , Option [ 'f' ]- [ "prefetch" ]- (ReqArg (\s o -> o { preFetch = read s }) "INT")- ("Prefetch count. (0=unlimited, 1=off, default: " ++- show (preFetch def) ++ ")")- , Option [ 's' ]- [ "heartbeats" ]- (ReqArg (\s o -> o { heartBeat = (Just (read s)) }) "INT")- "heartbeat interval (0=disable, default: set by server)"- , Option [ 't' ]- [ "tempdir", "target" ]- (OptArg (\s o -> o { tempDir = Just (fromMaybe "/tmp" s) })- "DIR")- "tempdir (default: no file creation, -t without arg: /tmp)"- , Option [ 'a' ]- [ "args" ]- (ReqArg (\s o -> o { additionalArgs = s : (additionalArgs o)- })- "ARG")- "additional argument for -X callback"- , Option [ 'n' ]- [ "name" ]- (ReqArg (\s o -> o { connectionName = Just s }) "NAME")- "connection name, will be shown in RabbitMQ web interface"- , Option [ 'l' ]- [ "charlimit" ]- (ReqArg (\s o -> o { anRiss = Just (read s :: Int) }) "INT")- "limit number of shown body chars (default: unlimited)"- ]--usage :: String-usage = "\n\- \amqp-utils " ++- (showVersion version) ++- "\n\n\- \Usage:\n\- \konsum [options] [bindingkey] [body byte count to show]\n\n\- \Options:"---- apply options onto argstring-parseargs :: [String] -> IO Args-parseargs argstring = case getOpt Permute options argstring of- (o, [], []) -> return $ foldl (flip id) def o- (_, _, errs) -> ioError $ userError $ concat errs ++ usageInfo usage options+ (q, _, _) <-+ declareQueue+ chan+ newQueue {queueExclusive = True, queueName = T.pack tmpqname}+ mapM_+ (\(xchange, bkey) ->+ bindQueue chan q (T.pack xchange) (T.pack bkey) >>+ printparam' "binding" (xchange ++ ":" ++ bkey))+ (if null bindlist+ then [(x, "#")]+ else bindlist)+ return q --- process received message-myCallback :: Maybe Int- -> Maybe String- -> Maybe String- -> [String]- -> ThreadId- -> (Message, Envelope)- -> IO ()+-- | process received message+myCallback ::+ Maybe Int+ -> Maybe String+ -> Maybe String+ -> [String]+ -> ThreadId+ -> (Message, Envelope)+ -> IO () myCallback anR filePr tempD addi tid m@(_, envi) = do- let numstring = show $ envDeliveryTag envi- hr $ "BEGIN " ++ numstring- now <- getZonedTime- callbackoptions <- printmsg m anR now- either (\e -> printparam' "ERROR" (show (e :: X.SomeException)) >>- rejectEnv envi True)- return- =<< X.try (optionalFileStuff m- callbackoptions- addi- numstring- tempD- filePr- tid)- hr $ "END " ++ numstring+ let numstring = show $ envDeliveryTag envi+ hr $ "BEGIN " ++ numstring+ now <- getZonedTime+ callbackoptions <- printmsg m anR now+ either+ (\e ->+ printparam' "ERROR" (show (e :: X.SomeException)) >> rejectEnv envi True)+ return =<<+ X.try (optionalFileStuff m callbackoptions addi numstring tempD filePr tid)+ hr $ "END " ++ numstring --- if the message is to be saved+-- | if the message is to be saved -- and maybe processed further-optionalFileStuff :: (Message, Envelope)- -> [String]- -> [String]- -> String- -> Maybe String- -> Maybe String- -> ThreadId- -> IO ()+optionalFileStuff ::+ (Message, Envelope)+ -> [String]+ -> [String]+ -> String+ -> Maybe String+ -> Maybe String+ -> ThreadId+ -> IO () optionalFileStuff (msg, envi) callbackoptions addi numstring tempD filePr tid = do- path <- saveFile tempD numstring (msgBody msg)- printparam "saved to" path- let callbackcmdline = liftM2 (constructCallbackCmdLine callbackoptions- addi- numstring)- filePr- path- printparam "calling" $ fmap unwords callbackcmdline- maybe (ackEnv envi)- (\c -> forkFinally (doProc numstring envi c)- (either (throwTo tid) return) >>- return ())- callbackcmdline+ path <- saveFile tempD numstring (msgBody msg)+ printparam "saved to" path+ let callbackcmdline =+ liftM2+ (constructCallbackCmdLine callbackoptions addi numstring)+ filePr+ path+ printparam "calling" $ fmap unwords callbackcmdline+ maybe+ (ackEnv envi)+ (\c ->+ forkFinally (doProc numstring envi c) (either (throwTo tid) return) >>+ return ())+ callbackcmdline --- save message into temp file+-- | save message into temp file saveFile :: Maybe String -> String -> BL.ByteString -> IO (Maybe String) saveFile Nothing _ _ = return Nothing saveFile (Just tempD) numstring body = do- (p, h) <- openBinaryTempFileWithDefaultPermissions tempD- ("konsum-" ++- numstring ++ "-.tmp")- BL.hPut h body- hClose h- return $ Just p+ (p, h) <-+ openBinaryTempFileWithDefaultPermissions+ tempD+ ("konsum-" ++ numstring ++ "-.tmp")+ BL.hPut h body+ hClose h+ return $ Just p --- construct cmdline for callback script-constructCallbackCmdLine :: [String]- -> [String]- -> String- -> String- -> String- -> [String]+-- | construct cmdline for callback script+constructCallbackCmdLine ::+ [String] -> [String] -> String -> String -> String -> [String] constructCallbackCmdLine opts addi num exe path =- exe : "-f" : path : "-n" : num : opts ++ addi+ exe : "-f" : path : "-n" : num : opts ++ addi --- call callback script+-- | call callback script doProc :: String -> Envelope -> [String] -> IO ()-doProc numstring envi (exe : args) = do- (_, _, _, processhandle) <- createProcess (proc exe args) { std_out = Inherit- , std_err = Inherit- }- exitcode <- waitForProcess processhandle- printparam' (numstring ++ " call returned") $ show exitcode- case exitcode of- ExitSuccess -> ackEnv envi- ExitFailure _ -> rejectEnv envi True+doProc numstring envi (exe:args) = do+ (_, _, _, processhandle) <-+ createProcess (proc exe args) {std_out = Inherit, std_err = Inherit}+ exitcode <- waitForProcess processhandle+ printparam' (numstring ++ " call returned") $ show exitcode+ case exitcode of+ ExitSuccess -> ackEnv envi+ ExitFailure _ -> rejectEnv envi True doProc _ _ _ = return ()---- log marker-hr :: String -> IO ()-hr x = putStrLn hr' >> hFlush stdout- where- hr' = take 72 $ (take 25 hr'') ++ " " ++ x ++ " " ++ hr''- hr'' = repeat '-'--formatheaders :: ((T.Text, FieldValue) -> [a]) -> FieldTable -> [a]-formatheaders f (FieldTable ll) =- concat $ map f $ M.toList ll---- log formatting-fieldshow :: (T.Text, FieldValue) -> String-fieldshow (k, v) = "\n " ++ T.unpack k ++ ": " ++ valueshow v---- callback cmdline formatting-fieldshow' :: (T.Text, FieldValue) -> [String]-fieldshow' (k, v) = [ "-h", T.unpack k ++ "=" ++ valueshow v ]---- showing a FieldValue-valueshow :: FieldValue -> String-valueshow (FVString value) =- T.unpack value-valueshow (FVInt32 value) =- show value-valueshow value = show value---- skip showing body head if binary type-isimage :: Maybe String -> Bool-isimage Nothing = False-isimage (Just ctype)- | isPrefixOf "application/xml" ctype =- False- | isPrefixOf "application/json" ctype =- False- | otherwise = any (flip isPrefixOf ctype) [ "application", "image" ]---- show the first bytes of message body-anriss' :: Maybe Int -> BL.ByteString -> BL.ByteString-anriss' x = case x of- Nothing -> id- Just y -> BL.take (fromIntegral y)---- callback cmdline with optional parameters-printopt :: (String, Maybe String) -> [String]-printopt (_, Nothing) = []-printopt (opt, Just s) =- [ opt, s ]---- prints header and head on STDOUT and returns cmdline options to callback-printmsg :: (Message, Envelope) -> Maybe Int -> ZonedTime -> IO [String]-printmsg (msg, envi) anR now = do- mapM_ (uncurry printparam)- [ ("routing key", rkey)- , ("message-id", messageid)- , ("headers", headers)- , ("content-type", contenttype)- , ("content-encoding", contentencoding)- , ("redelivered", redeliv)- , ("timestamp", timestamp'')- , ("time now", now')- , ("size", size)- , ("priority", prio)- , ("type", mtype)- , ("user id", muserid)- , ("application id", mappid)- , ("cluster id", mclusterid)- , ("reply to", mreplyto)- , ("correlation id", mcorrid)- , ("expiration", mexp)- , ("delivery mode", mdelivmode)- ]- printbody (label, anriss)- return $- concat (map printopt- [ ("-r", rkey)- , ("-m", contenttype)- , ("-e", contentencoding)- , ("-i", messageid)- , ("-t", timestamp)- , ("-p", prio)- ] ++- maybeToList headers')- where- headers = fmap (formatheaders fieldshow) $ msgHeaders msg- headers' = fmap (formatheaders fieldshow') $ msgHeaders msg- body = msgBody msg- anriss = if isimage contenttype- then Nothing- else Just (anriss' anR body) :: Maybe BL.ByteString- anriss'' = maybe "" (\a -> "first " ++ (show a) ++ " bytes of ") anR- label = anriss'' ++ "body"- contenttype = fmap T.unpack $ msgContentType msg- contentencoding = fmap T.unpack $ msgContentEncoding msg- rkey = Just . T.unpack $ envRoutingKey envi- messageid = fmap T.unpack $ msgID msg- prio = fmap show $ msgPriority msg- mtype = fmap show $ msgType msg- muserid = fmap show $ msgUserID msg- mappid = fmap show $ msgApplicationID msg- mclusterid = fmap show $ msgClusterID msg- mreplyto = fmap show $ msgReplyTo msg- mcorrid = fmap show $ msgCorrelationID msg- mexp = fmap show $ msgExpiration msg- mdelivmode = fmap show $ msgDeliveryMode msg- size = Just . show $ BL.length body- redeliv = if envRedelivered envi then Just "YES" else Nothing- tz = zonedTimeZone now- nowutc = zonedTimeToUTCFLoor now- msgtime = msgTimestamp msg- msgtimeutc = fmap (posixSecondsToUTCTime . realToFrac) msgtime- timestamp = fmap show msgtime- timediff = fmap (difftime nowutc) msgtimeutc- now' = case timediff of- Just "now" -> Nothing- _ -> showtime tz $ Just nowutc- timestamp' = showtime tz msgtimeutc- timestamp'' = liftM3 (\a b c -> a ++ " (" ++ b ++ ") (" ++ c ++ ")")- timestamp- timestamp'- timediff--zonedTimeToUTCFLoor :: ZonedTime -> UTCTime-zonedTimeToUTCFLoor x = posixSecondsToUTCTime $- realToFrac ((floor .- utcTimeToPOSIXSeconds .- zonedTimeToUTC) x :: Timestamp)--showtime :: TimeZone -> Maybe UTCTime -> Maybe String-showtime tz = fmap (show . (utcToZonedTime tz))--difftime :: UTCTime -> UTCTime -> String-difftime now msg- | now == msg = "now"- | now > msg = diff ++ " ago"- | otherwise = diff ++ " in the future"- where- diff = show (diffUTCTime now msg)