opentelemetry-extra 0.4.0 → 0.4.1
raw patch · 6 files changed
+135/−60 lines, 6 filesdep +optparse-applicativedep +processdep ~basenew-component:exe:eventlog-to-tracyPVP: major bump suggested
API removals or changes: PVP suggests a major version bump
Dependencies added: optparse-applicative, process
Dependency ranges changed: base
API changes (from Hackage documentation)
+ OpenTelemetry.EventlogStreaming_Internal: EventLogFilename :: FilePath -> EventSource
+ OpenTelemetry.EventlogStreaming_Internal: EventLogHandle :: Handle -> WatDoOnEOF -> EventSource
+ OpenTelemetry.EventlogStreaming_Internal: data EventSource
- OpenTelemetry.EventlogStreaming_Internal: work :: WatDoOnEOF -> Timestamp -> Exporter Span -> Handle -> IO ()
+ OpenTelemetry.EventlogStreaming_Internal: work :: Timestamp -> Exporter Span -> EventSource -> IO ()
Files
- exe/eventlog-to-chrome/Main.hs +25/−5
- exe/eventlog-to-tracy/Main.hs +25/−0
- exe/eventlog-to-zipkin/Main.hs +4/−2
- opentelemetry-extra.cabal +10/−1
- src/OpenTelemetry/ChromeExporter.hs +1/−2
- src/OpenTelemetry/EventlogStreaming_Internal.hs +70/−50
exe/eventlog-to-chrome/Main.hs view
@@ -5,20 +5,40 @@ import OpenTelemetry.ChromeExporter import OpenTelemetry.EventlogStreaming_Internal import OpenTelemetry.Exporter+import Options.Applicative import System.Clock-import System.Environment (getArgs) import System.IO import Text.Printf +data ConsoleOptions = ConsoleOptions Command deriving (Show)++data Command = EventlogToChromeCmd FilePath deriving (Show)+ main :: IO () main = do- args <- getArgs- case args of- ["read", path] -> do+ (ConsoleOptions cmd) <- parseConsoleOptions+ case cmd of+ EventlogToChromeCmd path -> do let target_path = (path <> ".trace.json") printf "Converting %s to %s...\n" path target_path exporter <- createChromeSpanExporter target_path origin_timestamp <- fromIntegral . toNanoSecs <$> getTime Realtime- withFile path ReadMode (work StopOnEOF origin_timestamp exporter)+ work origin_timestamp exporter $ EventLogFilename path shutdown exporter putStrLn "\nAll done."++readEventlogFileCmdParser :: Parser Command+readEventlogFileCmdParser+ = EventlogToChromeCmd <$> argument str (metavar "FILE")++consoleOptionParser :: Parser ConsoleOptions+consoleOptionParser+ = ConsoleOptions+ <$> hsubparser+ (command "read"+ (info readEventlogFileCmdParser+ (progDesc "converts eventlog into chrome trace")))++parseConsoleOptions :: IO ConsoleOptions+parseConsoleOptions+ = execParser $ info (consoleOptionParser <**> helper) fullDesc
+ exe/eventlog-to-tracy/Main.hs view
@@ -0,0 +1,25 @@+module Main where++import System.Environment (getArgs)+import System.Exit+import System.Process++help :: IO ()+help = do+ putStrLn "Converts eventlog to Tracy format and launches the viewer."+ putStrLn "Path to eventlog is expected."+ exitFailure++main :: IO ()+main = do+ args <- getArgs+ case args of+ ["-h"] -> help+ ["--help"] -> help+ [eventlogFile] -> do+ let chromeFile = eventlogFile ++ ".trace.json"+ let tracyFile = eventlogFile ++ ".tracy"+ callProcess "eventlog-to-chrome" ["read", eventlogFile]+ callProcess "import-chrome" [chromeFile, tracyFile]+ callProcess "Tracy" [tracyFile]+ _ -> help
exe/eventlog-to-zipkin/Main.hs view
@@ -25,7 +25,7 @@ let service_name = T.pack $ takeBaseName path exporter <- createZipkinSpanExporter $ localhostZipkinConfig service_name origin_timestamp <- fromIntegral . toNanoSecs <$> getTime Realtime- withFile path ReadMode (work StopOnEOF origin_timestamp exporter)+ work origin_timestamp exporter $ EventLogFilename path shutdown exporter putStrLn "\nAll done." ("run" : program : "--" : args') -> do@@ -36,7 +36,9 @@ env <- (("GHCRTS", "-l -ol" <> pipe) :) <$> getEnvironment -- TODO(divanov): please append to existing GHCRTS instead of overwriting p <- startProcess (proc program args' & setEnv env) origin_timestamp <- fromIntegral . toNanoSecs <$> getTime Realtime- restreamer <- async $ withFile pipe ReadMode (work SleepAndRetryOnEOF origin_timestamp exporter)+ restreamer <- async $+ withFile pipe ReadMode (\input ->+ work origin_timestamp exporter $ EventLogHandle input SleepAndRetryOnEOF) waitExitCode p wait restreamer shutdown exporter
opentelemetry-extra.cabal view
@@ -2,7 +2,7 @@ name: opentelemetry-extra description: The OpenTelemetry Haskell Client https://opentelemetry.io category: OpenTelemetry-version: 0.4.0+version: 0.4.1 license-file: LICENSE license: Apache-2.0 author: Dmitry Ivanov@@ -130,3 +130,12 @@ clock, opentelemetry >= 0.4.0, opentelemetry-extra,+ optparse-applicative,++executable eventlog-to-tracy+ import: options+ main-is: Main.hs+ hs-source-dirs: exe/eventlog-to-tracy+ build-depends:+ base,+ process,
src/OpenTelemetry/ChromeExporter.hs view
@@ -4,7 +4,6 @@ import Data.Aeson import qualified Data.ByteString.Lazy as LBS-import qualified Data.HashMap.Strict as HM import OpenTelemetry.Common import OpenTelemetry.Exporter import System.IO@@ -45,7 +44,7 @@ createChromeSpanExporter :: FilePath -> IO (Exporter Span) createChromeSpanExporter path = do f <- openFile path WriteMode- hPutStrLn f "["+ hPutStrLn f "[ " pure $! Exporter ( \sps -> do
src/OpenTelemetry/EventlogStreaming_Internal.hs view
@@ -25,57 +25,77 @@ instance Hashable SpanId -work :: WatDoOnEOF -> Timestamp -> Exporter Span -> Handle -> IO ()-work wat_do_on_eof origin_timestamp exporter input = do+data EventSource+ = EventLogHandle Handle WatDoOnEOF+ | EventLogFilename FilePath++work :: Timestamp -> Exporter Span -> EventSource -> IO ()+work origin_timestamp exporter source = do d_ "Starting the eventlog reader" smgen <- R.initSMGen -- TODO(divanov): seed the random generator with something more random than current time- go (initialState origin_timestamp smgen) decodeEventLog+ let state0 = initialState origin_timestamp smgen+ case source of+ EventLogFilename path -> do+ readEventLogFromFile path >>= \case+ Right (dat -> Data {events}) -> do+ let go s [] = pure ()+ go s (e : es) = do+ dd_ "event" (evTime e, evCap e, evSpec e)+ case processEvent e s of+ (s', sps) -> do+ mapM_ (d_ . ("emit " <>) . show) sps+ _ <- export exporter sps+ go s' es+ go state0 $ sortEvents events+ Left err -> do+ putStrLn err+ EventLogHandle input wat_do_on_eof -> do+ let go s (Produce event next) = do+ case evSpec event of+ Shutdown {} -> do+ d_ "Shutdown-like event detected"+ CapDelete {} -> do+ d_ "Shutdown-like event detected"+ CapsetDelete {} -> do+ d_ "Shutdown-like event detected"+ _ -> do+ -- d_ "go Produce"+ dd_ "event" (evTime event, evCap event, evSpec event)+ let (s', sps) = processEvent event s+ _ <- export exporter sps+ -- print s'+ mapM_ (d_ . ("emit " <>) . show) sps+ go s' next+ go s d@(Consume consume) = do+ -- d_ "go Consume"+ eof <- hIsEOF input+ case eof of+ False -> do+ chunk <- B.hGetSome input 4096+ -- printf "chunk = %d bytes\n" (B.length chunk)+ if B.null chunk+ then do+ -- d_ "chunk is null"+ threadDelay 1000 -- TODO(divanov): remove the sleep by replacing the hGetSome with something that blocks until data is available+ go s d+ else do+ -- d_ "chunk is not null"+ go s $ consume chunk+ True -> do+ d_ "EOF"+ case wat_do_on_eof of+ StopOnEOF -> pure ()+ SleepAndRetryOnEOF -> do+ threadDelay 1000+ go s d+ go _ (Done _) = do+ d_ "go Done"+ pure ()+ go _ (Error _leftover err) = do+ d_ "go Error"+ d_ err+ go state0 decodeEventLog d_ "no more work"- where- go s (Produce event next) = do- case evSpec event of- Shutdown {} -> do- d_ "Shutdown-like event detected"- CapDelete {} -> do- d_ "Shutdown-like event detected"- CapsetDelete {} -> do- d_ "Shutdown-like event detected"- _ -> do- -- d_ "go Produce"- dd_ "event" (evTime event, evCap event, evSpec event)- let (s', sps) = processEvent event s- _ <- export exporter sps- -- print s'- mapM_ (d_ . ("emit " <>) . show) sps- go s' next- go s d@(Consume consume) = do- -- d_ "go Consume"- eof <- hIsEOF input- case eof of- False -> do- chunk <- B.hGetSome input 4096- -- printf "chunk = %d bytes\n" (B.length chunk)- if B.null chunk- then do- -- d_ "chunk is null"- threadDelay 1000 -- TODO(divanov): remove the sleep by replacing the hGetSome with something that blocks until data is available- go s d- else do- -- d_ "chunk is not null"- go s $ consume chunk- True -> do- d_ "EOF"- case wat_do_on_eof of- StopOnEOF -> pure ()- SleepAndRetryOnEOF -> do- threadDelay 1000- go s d- go _ (Done _) = do- d_ "go Done"- pure ()- go _ (Error _leftover err) = do- d_ "go Error"- d_ err data State = S { originTimestamp :: Timestamp,@@ -186,13 +206,13 @@ ("ot2" : "set" : "tag" : serial_text : k : v) -> let serial = read (T.unpack serial_text) in case HM.lookup serial serial2sid of- Nothing -> error $ "set tag: span id not found for serial" <> T.unpack serial_text+ Nothing -> error $ "set tag: span id not found for serial " <> T.unpack serial_text Just span_id -> (modifySpan span_id (setTag k (T.unwords v)) st, []) ["ot2", "set", "traceid", serial_text, trace_id_text] -> let serial = read (T.unpack serial_text) trace_id = TId (read ("0x" <> T.unpack trace_id_text)) in case HM.lookup serial serial2sid of- Nothing -> error $ "set traceid: span id not found for serial" <> T.unpack serial_text+ Nothing -> error $ "set traceid: span id not found for serial " <> T.unpack serial_text Just span_id -> ( (modifySpan span_id (setTraceId trace_id) st) { traceMap = HM.insert tid trace_id traceMap