packages feed

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 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