packages feed

opentelemetry-extra 0.3.2 → 0.4.0

raw patch · 9 files changed

+429/−280 lines, 9 filesdep +QuickCheckdep +text-showdep +typed-processdep ~opentelemetryPVP ok

version bump matches the API change (PVP)

Dependencies added: QuickCheck, text-show, typed-process

Dependency ranges changed: opentelemetry

API changes (from Hackage documentation)

- OpenTelemetry.EventlogStreaming_Internal: [spanStacks] :: State -> HashMap ThreadId (NonEmpty Span)
- OpenTelemetry.EventlogStreaming_Internal: [specificSpans] :: State -> HashMap Word64 Span
- OpenTelemetry.EventlogStreaming_Internal: beginSpecificSpan :: TraceId -> SpanId -> Text -> Timestamp -> State -> (State, [Span])
- OpenTelemetry.EventlogStreaming_Internal: endSpecificSpan :: SpanId -> Timestamp -> State -> (State, [Span])
- OpenTelemetry.EventlogStreaming_Internal: modifyAllSpans :: (Span -> Span) -> State -> State
- OpenTelemetry.EventlogStreaming_Internal: popSpan :: HasCallStack => ThreadId -> Timestamp -> State -> (State, [Span])
- OpenTelemetry.EventlogStreaming_Internal: popSpansAcrossAllThreads :: HasCallStack => Timestamp -> State -> (State, [Span])
- OpenTelemetry.EventlogStreaming_Internal: pushGCSpans :: HasCallStack => State -> Timestamp -> State
- OpenTelemetry.EventlogStreaming_Internal: pushSpan :: HasCallStack => ThreadId -> Text -> Timestamp -> State -> State
- OpenTelemetry.Tracer: Tracer :: !HashMap threadId (NonEmpty Span) -> Tracer threadId
- OpenTelemetry.Tracer: [tracerSpanStacks] :: Tracer threadId -> !HashMap threadId (NonEmpty Span)
- OpenTelemetry.Tracer: createTracer :: (Hashable tid, Eq tid) => IO (Tracer tid)
- OpenTelemetry.Tracer: data Tracer threadId
- OpenTelemetry.Tracer: instance GHC.Classes.Eq threadId => GHC.Classes.Eq (OpenTelemetry.Tracer.Tracer threadId)
- OpenTelemetry.Tracer: instance GHC.Show.Show threadId => GHC.Show.Show (OpenTelemetry.Tracer.Tracer threadId)
- OpenTelemetry.Tracer: tracerGetCurrentActiveSpan :: (Hashable tid, Eq tid) => Tracer tid -> tid -> Maybe Span
- OpenTelemetry.Tracer: tracerPopSpan :: (Eq tid, Hashable tid) => Tracer tid -> tid -> (Maybe Span, Tracer tid)
- OpenTelemetry.Tracer: tracerPushSpan :: (Eq tid, Hashable tid) => Tracer tid -> tid -> Span -> Tracer tid
+ OpenTelemetry.Common: [spanThreadId] :: Span -> Word32
+ OpenTelemetry.EventlogStreaming_Internal: SleepAndRetryOnEOF :: WatDoOnEOF
+ OpenTelemetry.EventlogStreaming_Internal: StopOnEOF :: WatDoOnEOF
+ OpenTelemetry.EventlogStreaming_Internal: [serial2sid] :: State -> HashMap Word64 SpanId
+ OpenTelemetry.EventlogStreaming_Internal: [spans] :: State -> HashMap SpanId Span
+ OpenTelemetry.EventlogStreaming_Internal: [thread2sid] :: State -> HashMap ThreadId SpanId
+ OpenTelemetry.EventlogStreaming_Internal: data WatDoOnEOF
+ OpenTelemetry.EventlogStreaming_Internal: emitSpan :: Word64 -> SpanId -> State -> (State, Span)
+ OpenTelemetry.EventlogStreaming_Internal: instance Data.Hashable.Class.Hashable OpenTelemetry.SpanContext.SpanId
+ OpenTelemetry.EventlogStreaming_Internal: inventSpanId :: Word64 -> State -> (State, SpanId)
- OpenTelemetry.Common: Span :: {-# UNPACK #-} !SpanContext -> Text -> !Timestamp -> !Timestamp -> !HashMap Text TagValue -> [SpanEvent] -> !SpanStatus -> Maybe SpanId -> Span
+ OpenTelemetry.Common: Span :: {-# UNPACK #-} !SpanContext -> Text -> Word32 -> !Timestamp -> !Timestamp -> !HashMap Text TagValue -> [SpanEvent] -> !SpanStatus -> Maybe SpanId -> Span
- OpenTelemetry.EventlogStreaming_Internal: S :: Timestamp -> IntMap ThreadId -> HashMap ThreadId (NonEmpty Span) -> HashMap ThreadId TraceId -> HashMap Word64 Span -> SMGen -> State
+ OpenTelemetry.EventlogStreaming_Internal: S :: Timestamp -> IntMap ThreadId -> HashMap SpanId Span -> HashMap ThreadId TraceId -> HashMap Word64 SpanId -> HashMap ThreadId SpanId -> SMGen -> State
- OpenTelemetry.EventlogStreaming_Internal: modifySpan :: HasCallStack => ThreadId -> (Span -> Span) -> State -> State
+ OpenTelemetry.EventlogStreaming_Internal: modifySpan :: HasCallStack => SpanId -> (Span -> Span) -> State -> State
- OpenTelemetry.EventlogStreaming_Internal: work :: Timestamp -> Exporter Span -> Handle -> IO ()
+ OpenTelemetry.EventlogStreaming_Internal: work :: WatDoOnEOF -> Timestamp -> Exporter Span -> Handle -> IO ()

Files

exe/eventlog-to-chrome/Main.hs view
@@ -19,6 +19,6 @@       printf "Converting %s to %s...\n" path target_path       exporter <- createChromeSpanExporter target_path       origin_timestamp <- fromIntegral . toNanoSecs <$> getTime Realtime-      withFile path ReadMode (work origin_timestamp exporter)+      withFile path ReadMode (work StopOnEOF origin_timestamp exporter)       shutdown exporter       putStrLn "\nAll done."
exe/eventlog-to-zipkin/Main.hs view
@@ -2,14 +2,18 @@  module Main where +import Control.Concurrent.Async+import Control.Monad+import Data.Function ((&)) import qualified Data.Text as T import OpenTelemetry.EventlogStreaming_Internal import OpenTelemetry.Exporter import OpenTelemetry.ZipkinExporter import System.Clock-import System.Environment (getArgs)+import System.Environment (getArgs, getEnvironment) import System.FilePath import System.IO+import System.Process.Typed import Text.Printf  main :: IO ()@@ -21,7 +25,20 @@       let service_name = T.pack $ takeBaseName path       exporter <- createZipkinSpanExporter $ localhostZipkinConfig service_name       origin_timestamp <- fromIntegral . toNanoSecs <$> getTime Realtime-      withFile path ReadMode (work origin_timestamp exporter)+      withFile path ReadMode (work StopOnEOF origin_timestamp exporter)+      shutdown exporter+      putStrLn "\nAll done."+    ("run" : program : "--" : args') -> do+      printf "Streaming eventlog of %s to Zipkin...\n" program+      exporter <- createZipkinSpanExporter $ localhostZipkinConfig (T.pack program)+      let pipe = program <> "-opentelemetry.pipe"+      runProcess $ proc "mkfifo" [pipe]+      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)+      waitExitCode p+      wait restreamer       shutdown exporter       putStrLn "\nAll done."     _ -> do
opentelemetry-extra.cabal view
@@ -2,7 +2,7 @@ name:                opentelemetry-extra description:         The OpenTelemetry Haskell Client https://opentelemetry.io category:            OpenTelemetry-version: 0.3.2+version: 0.4.0 license-file:        LICENSE license:             Apache-2.0 author:              Dmitry Ivanov@@ -62,9 +62,10 @@     http-client,     http-client-tls,     http-types,-    opentelemetry >= 0.3.0,+    opentelemetry >= 0.4.0,     random >= 1.1,     scientific,+    text-show,     splitmix,     stm,     text,@@ -73,7 +74,6 @@   hs-source-dirs: src   exposed-modules:     OpenTelemetry.Common-    OpenTelemetry.Tracer     OpenTelemetry.Debug     OpenTelemetry.ChromeExporter     OpenTelemetry.EventlogStreaming_Internal@@ -84,18 +84,24 @@   type:       exitcode-stdio-1.0   main-is:    Main.hs   other-modules:-    TestCommon     TestPropagation+    TestEventlogStreaming   hs-source-dirs: unit-tests   build-depends:     base,+    QuickCheck,     bytestring,+    ghc-events,+    hashable,+    opentelemetry >= 0.4.0,+    opentelemetry-extra,     tasty,-    tasty-quickcheck,-    tasty-hunit,     tasty-discover,-    opentelemetry >= 0.3.0,-    opentelemetry-extra+    tasty-hunit,+    tasty-quickcheck,+    text,+    text-show,+    unordered-containers  executable eventlog-to-zipkin   import: options@@ -109,9 +115,10 @@     filepath,     http-client,     http-client-tls,-    opentelemetry >= 0.3.0,+    opentelemetry >= 0.4.0,     opentelemetry-extra,-    text+    text,+    typed-process  executable eventlog-to-chrome   import: options@@ -121,5 +128,5 @@     base,     exceptions,     clock,-    opentelemetry >= 0.3.0,+    opentelemetry >= 0.4.0,     opentelemetry-extra,
src/OpenTelemetry/ChromeExporter.hs view
@@ -5,12 +5,9 @@ import Data.Aeson import qualified Data.ByteString.Lazy as LBS import qualified Data.HashMap.Strict as HM-import qualified Data.Text as T import OpenTelemetry.Common import OpenTelemetry.Exporter-import OpenTelemetry.SpanContext import System.IO-import Text.Read  newtype ChromeBeginSpan = ChromeBegin Span @@ -26,28 +23,22 @@  instance ToJSON ChromeBeginSpan where   toJSON (ChromeBegin Span {..}) =-    let threadId = case HM.lookup "tid" spanTags of-          Just (IntTagValue t) -> t-          _ -> 1-     in object-          [ "ph" .= ("B" :: String),-            "name" .= spanOperation,-            "pid" .= (1 :: Int),-            "tid" .= threadId,-            "ts" .= (div spanStartedAt 1000),-            "args" .= fmap ChromeTagValue spanTags-          ]+    object+      [ "ph" .= ("B" :: String),+        "name" .= spanOperation,+        "pid" .= (1 :: Int),+        "tid" .= spanThreadId,+        "ts" .= (div spanStartedAt 1000),+        "args" .= fmap ChromeTagValue spanTags+      ]  instance ToJSON ChromeEndSpan where   toJSON (ChromeEnd Span {..}) =-    let threadId = case HM.lookup "tid" spanTags of-          Just (IntTagValue t) -> t-          _ -> 1-     in object+     object           [ "ph" .= ("E" :: String),             "name" .= spanOperation,             "pid" .= (1 :: Int),-            "tid" .= threadId,+            "tid" .= spanThreadId,             "ts" .= (div spanFinishedAt 1000)           ] 
src/OpenTelemetry/Common.hs view
@@ -33,17 +33,17 @@ instance ToTagValue Int where   toTagValue = IntTagValue -data Span-  = Span-      { spanContext :: {-# UNPACK #-} !SpanContext,-        spanOperation :: T.Text,-        spanStartedAt :: !Timestamp,-        spanFinishedAt :: !Timestamp,-        spanTags :: !(HM.HashMap T.Text TagValue),-        spanEvents :: [SpanEvent],-        spanStatus :: !SpanStatus,-        spanParentId :: Maybe SpanId-      }+data Span = Span+  { spanContext :: {-# UNPACK #-} !SpanContext,+    spanOperation :: T.Text,+    spanThreadId :: Word32,+    spanStartedAt :: !Timestamp,+    spanFinishedAt :: !Timestamp,+    spanTags :: !(HM.HashMap T.Text TagValue),+    spanEvents :: [SpanEvent],+    spanStatus :: !SpanStatus,+    spanParentId :: Maybe SpanId+  }   deriving (Show, Eq)  spanTraceId :: Span -> TraceId@@ -52,12 +52,11 @@ spanId :: Span -> SpanId spanId Span {spanContext = SpanContext sid _} = sid -data SpanEvent-  = SpanEvent-      { spanEventTimestamp :: !Timestamp,-        spanEventKey :: !T.Text,-        spanEventValue :: !T.Text-      }+data SpanEvent = SpanEvent+  { spanEventTimestamp :: !Timestamp,+    spanEventKey :: !T.Text,+    spanEventValue :: !T.Text+  }   deriving (Show, Eq)  data SpanStatus = OK@@ -67,16 +66,14 @@   = Event T.Text Timestamp   deriving (Show, Eq) -data SpanProcessor-  = SpanProcessor-      { onStart :: Span -> IO (),-        onEnd :: Span -> IO ()-      }+data SpanProcessor = SpanProcessor+  { onStart :: Span -> IO (),+    onEnd :: Span -> IO ()+  } -data OpenTelemetryConfig-  = OpenTelemetryConfig-      { otcSpanExporter :: Exporter Span-      }+data OpenTelemetryConfig = OpenTelemetryConfig+  { otcSpanExporter :: Exporter Span+  }  now64 :: IO Timestamp now64 = do
src/OpenTelemetry/EventlogStreaming_Internal.hs view
@@ -5,8 +5,8 @@ import Control.Concurrent (threadDelay) import qualified Data.ByteString as B import qualified Data.HashMap.Strict as HM+import Data.Hashable import qualified Data.IntMap as IM-import Data.List.NonEmpty as NE import Data.Maybe import qualified Data.Text as T import Data.Word@@ -14,7 +14,6 @@ import GHC.RTS.Events.Incremental import GHC.Stack import OpenTelemetry.Common hiding (Event, Timestamp)-import qualified OpenTelemetry.Common as OTel import OpenTelemetry.Debug import OpenTelemetry.Exporter import OpenTelemetry.SpanContext@@ -22,8 +21,12 @@ import qualified System.Random.SplitMix as R import Text.Printf -work :: Timestamp -> Exporter Span -> Handle -> IO ()-work origin_timestamp exporter input = do+data WatDoOnEOF = StopOnEOF | SleepAndRetryOnEOF++instance Hashable SpanId++work :: WatDoOnEOF -> Timestamp -> Exporter Span -> Handle -> IO ()+work wat_do_on_eof origin_timestamp exporter input = 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@@ -62,6 +65,11 @@               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 ()@@ -72,71 +80,189 @@ data State = S   { originTimestamp :: Timestamp,     threadMap :: IM.IntMap ThreadId,-    spanStacks :: HM.HashMap ThreadId (NonEmpty Span),+    spans :: HM.HashMap SpanId Span,     traceMap :: HM.HashMap ThreadId TraceId,-    specificSpans :: HM.HashMap Word64 Span,+    serial2sid :: HM.HashMap Word64 SpanId,+    thread2sid :: HM.HashMap ThreadId SpanId,     randomGen :: R.SMGen   }   deriving (Show)  initialState :: Word64 -> R.SMGen -> State-initialState timestamp = S timestamp mempty mempty mempty mempty+initialState timestamp = S timestamp mempty mempty mempty mempty mempty +inventSpanId :: Word64 -> State -> (State, SpanId)+inventSpanId serial st = (st {serial2sid = HM.insert serial sid (serial2sid st)}, sid)+  where+    sid = SId serial -- TODO: use random generator instead+ processEvent :: Event -> State -> (State, [Span]) processEvent (Event ts ev m_cap) st@(S {..}) =   let now = originTimestamp + ts       m_thread_id = m_cap >>= flip IM.lookup threadMap       m_trace_id = m_thread_id >>= flip HM.lookup traceMap    in case (ev, m_cap, m_thread_id) of-        (WallClockTime {sec, nsec}, _, _) -> (st {originTimestamp = sec * 1_000_000_000 + fromIntegral nsec - ts}, [])+        (WallClockTime {sec, nsec}, _, _) ->+          (st {originTimestamp = sec * 1_000_000_000 + fromIntegral nsec - ts}, [])         (CreateThread new_tid, _, _) ->-          case m_trace_id of-            Just trace_id -> (st {traceMap = HM.insert new_tid trace_id traceMap}, [])-            _ -> (st, [])+          let trace_id = case m_trace_id of+                Just t -> t+                Nothing -> TId originTimestamp -- TODO: something more random+           in (st {traceMap = HM.insert new_tid trace_id traceMap}, [])         (RunThread tid, Just cap, _) ->           (st {threadMap = IM.insert cap tid threadMap}, [])-        (StopThread _ tstatus, Just cap, _)-          | isTerminalThreadStatus tstatus -> (st {threadMap = IM.delete cap threadMap}, [])-        (StartGC, _, _) -> (pushGCSpans st now, [])-        (GCStatsGHC {gen}, _, _) -> (modifyAllSpans (setTag "gen" gen) st, [])-        (EndGC, _, _) -> popSpansAcrossAllThreads now st-        (HeapAllocated {allocBytes}, _, Just tid) ->-          (modifySpan tid (addEvent now "heap_alloc_bytes" (showT allocBytes)) st, [])+        (StopThread tid tstatus, Just cap, _)+          | isTerminalThreadStatus tstatus ->+            ( st+                { threadMap = IM.delete cap threadMap,+                  traceMap = HM.delete tid traceMap+                },+              []+            )+        -- (StartGC, _, _) ->+        --   (pushGCSpans st now, [])+        -- (GCStatsGHC {gen}, _, _) ->+        --   (modifyAllSpans (setTag "gen" gen) st, [])+        -- (EndGC, _, _) ->+        --   popSpansAcrossAllThreads now st+        -- (HeapAllocated {allocBytes}, _, Just tid) ->+        --   (modifySpan tid (addEvent now "heap_alloc_bytes" (showT allocBytes)) st, [])         (UserMessage {msg}, _, fromMaybe 1 -> tid) -> case T.words msg of-          ("ot1" : "begin" : "specific" : "span" : trace_id_text : span_id_text : name) ->-            let trace_id = TId (read ("0x" <> T.unpack trace_id_text))-                span_id = SId (read ("0x" <> T.unpack span_id_text))-             in beginSpecificSpan trace_id span_id (T.intercalate " " name) now st-          ("ot1" : "end" : "specific" : "span" : trace_id_text : span_id_text : _) ->-            let trace_id = TId (read ("0x" <> T.unpack trace_id_text))-                span_id = SId (read ("0x" <> T.unpack span_id_text))-             in endSpecificSpan span_id now st-          ("ot1" : "begin" : "span" : name) ->-            (pushSpan tid (T.intercalate " " name) now st, [])-          ("ot1" : "end" : "span" : _) -> popSpan tid now st-          ("ot1" : "set" : "tag" : k : v) -> (modifySpan tid (setTag k (T.unwords v)) st, [])-          ["ot1", "set", "traceid", trace_id_text] ->-            let trace_id = TId (read ("0x" <> T.unpack trace_id_text))-             in ( (modifySpan tid (setTraceId trace_id) st)-                    { traceMap = HM.insert tid trace_id traceMap-                    },-                  []-                )-          ["ot1", "set", "spanid", span_id] ->-            (modifySpan tid (setSpanId (SId (read ("0x" <> T.unpack span_id)))) st, [])-          ["ot1", "set", "parent", trace_id_text, span_id_text] ->+          ("ot2" : "begin" : "span" : serial_text : name) ->+            let serial = read (T.unpack serial_text)+                operation = T.intercalate " " name+             in case HM.lookup serial serial2sid of+                  Nothing ->+                    let (st', span_id) = inventSpanId serial st+                        parent = HM.lookup tid thread2sid+                        sp =+                          Span+                            { spanContext = SpanContext span_id (fromMaybe (TId 42) m_trace_id),+                              spanOperation = operation,+                              spanThreadId = tid,+                              spanStartedAt = now,+                              spanFinishedAt = 0,+                              spanTags = mempty,+                              spanEvents = mempty,+                              spanStatus = OK,+                              spanParentId = parent+                            }+                     in ( st'+                            { spans = HM.insert span_id sp spans,+                              thread2sid = HM.insert tid span_id thread2sid+                            },+                          []+                        )+                  Just span_id ->+                    let (st', sp) = emitSpan serial span_id st+                     in (st', [sp {spanOperation = operation, spanStartedAt = now, spanThreadId = tid}])+          ["ot2", "end", "span", serial_text] ->+            let serial = read (T.unpack serial_text)+             in case HM.lookup serial serial2sid of+                  Nothing ->+                    let (st', span_id) = inventSpanId serial st+                        parent = HM.lookup tid thread2sid+                        sp =+                          Span+                            { spanContext = SpanContext span_id (fromMaybe (TId 42) m_trace_id),+                              spanOperation = "",+                              spanThreadId = tid,+                              spanStartedAt = 0,+                              spanFinishedAt = now,+                              spanTags = mempty,+                              spanEvents = mempty,+                              spanStatus = OK,+                              spanParentId = parent+                            }+                     in ( st'+                            { spans = HM.insert span_id sp spans,+                              thread2sid = HM.insert tid span_id thread2sid+                            },+                          []+                        )+                  Just span_id ->+                    let (st', sp) = emitSpan serial span_id st+                     in (st', [sp {spanFinishedAt = now}])+          ("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+                  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+                  Just span_id ->+                    ( (modifySpan span_id (setTraceId trace_id) st)+                        { traceMap = HM.insert tid trace_id traceMap+                        },+                      []+                    )+          ["ot2", "set", "spanid", serial_text, new_span_id_text] ->+            let serial = read (T.unpack serial_text)+             in case HM.lookup serial serial2sid of+                  Just old_span_id -> (modifySpan old_span_id (setSpanId (SId (read ("0x" <> T.unpack new_span_id_text)))) st, [])+                  Nothing -> error $ "set spanid " <> T.unpack serial_text <> " " <> T.unpack new_span_id_text <> ": span id not found"+          ["ot2", "set", "parent", serial_text, trace_id_text, parent_span_id_text] ->             let trace_id = TId (read ("0x" <> T.unpack trace_id_text))-                sid = SId (read ("0x" <> T.unpack span_id_text))-             in ( (modifySpan tid (setParent trace_id sid) st)-                    { traceMap = HM.insert tid trace_id traceMap-                    },-                  []-                )-          ("ot1" : "add" : "event" : k : v) -> (modifySpan tid (addEvent now k (T.unwords v)) st, [])-          ("ot1" : rest) -> error $ printf "Unrecognized %s" (show rest)+                serial = read (T.unpack serial_text)+                psid = SId (read ("0x" <> T.unpack parent_span_id_text))+             in case HM.lookup serial serial2sid of+                  Just span_id ->+                    ( (modifySpan span_id (setParent trace_id psid) st)+                        { traceMap = HM.insert tid trace_id traceMap+                        },+                      []+                    )+                  Nothing -> error $ "set parent: span not found for serial " <> show serial+          ("ot2" : "add" : "event" : serial_text : k : v) ->+            let serial = read (T.unpack serial_text)+             in case HM.lookup serial serial2sid of+                  Just span_id -> (modifySpan span_id (addEvent now k (T.unwords v)) st, [])+                  Nothing -> error $ "add event: span not found for serial " <> show serial+          ("ot2" : rest) -> error $ printf "Unrecognized %s" (show rest)           _ -> (st, [])         _ -> (st, []) +-- beginSpan :: TraceId -> SpanId -> T.Text -> OTel.Timestamp -> State -> (State, [Span])+-- beginSpan trace_id span_id@(SId s) name timestamp st =+--   case HM.lookup s (specificSpans st) of+--     Just sp -> (st {specificSpans = HM.delete s (specificSpans st)}, [sp {spanStartedAt = timestamp, spanOperation = name, spanContext = SpanContext span_id trace_id}])+--     Nothing ->+--       (st {specificSpans = HM.insert s sp (specificSpans st)}, [])+--       where+--         sp =+--           Span+--             { spanContext = SpanContext span_id trace_id,+--               spanOperation = name,+--               spanStartedAt = timestamp,+--               spanFinishedAt = 0,+--               spanTags = mempty,+--               spanEvents = mempty,+--               spanStatus = OK,+--               spanParentId = Nothing+--             }++-- endSpan :: Word64 -> OTel.Timestamp -> State -> (State, [Span])+-- endSpan serial timestamp st =+--   case HM.lookup serial (serial2sid st) of+--     Just span_id@(SId s) -> (st {specificSpans = HM.delete s (specificSpans st)}, [sp {spanFinishedAt = timestamp}])+--     Nothing ->+--       (st {specificSpans = HM.insert s sp (specificSpans st)}, [])+--       where+--         sp =+--           Span+--             { spanContext = SpanContext span_id (TId 0),+--               spanOperation = "unknown",+--               spanStartedAt = 0,+--               spanFinishedAt = timestamp,+--               spanTags = mempty,+--               spanEvents = mempty,+--               spanStatus = OK,+--               spanParentId = Nothing+--             }+ setTag :: ToTagValue v => T.Text -> v -> Span -> Span setTag k v sp =   sp@@ -168,105 +294,63 @@     new_events = ev : spanEvents sp     ev = SpanEvent ts k v -modifyAllSpans :: (Span -> Span) -> State -> State-modifyAllSpans f st =-  st-    { spanStacks =-        fmap-          (\(sp :| sps) -> (f sp :| sps))-          (spanStacks st)-    }--modifySpan :: HasCallStack => ThreadId -> (Span -> Span) -> State -> State-modifySpan tid f st =-  st-    { spanStacks =-        HM.update (\(sp :| sps) -> Just (f sp :| sps)) tid (spanStacks st)-    }--beginSpecificSpan :: TraceId -> SpanId -> T.Text -> OTel.Timestamp -> State -> (State, [Span])-beginSpecificSpan trace_id span_id@(SId s) name timestamp st =-  case HM.lookup s (specificSpans st) of-    Just sp -> (st {specificSpans = HM.delete s (specificSpans st)}, [sp {spanStartedAt = timestamp, spanOperation = name, spanContext = SpanContext span_id trace_id}])-    Nothing ->-      (st {specificSpans = HM.insert s sp (specificSpans st)}, [])-      where-        sp =-          Span-            { spanContext = SpanContext span_id trace_id,-              spanOperation = name,-              spanStartedAt = timestamp,-              spanFinishedAt = 0,-              spanTags = mempty,-              spanEvents = mempty,-              spanStatus = OK,-              spanParentId = Nothing-            }+-- modifyAllSpans :: (Span -> Span) -> State -> State+-- modifyAllSpans f st =+--   st+--     { spanStacks =+--         fmap+--           (\(sp :| sps) -> (f sp :| sps))+--           (spanStacks st)+--     } -endSpecificSpan :: SpanId -> OTel.Timestamp -> State -> (State, [Span])-endSpecificSpan span_id@(SId s) timestamp st =-  case HM.lookup s (specificSpans st) of-    Just sp -> (st {specificSpans = HM.delete s (specificSpans st)}, [sp {spanFinishedAt = timestamp}])-    Nothing ->-      (st {specificSpans = HM.insert s sp (specificSpans st)}, [])-      where-        sp =-          Span-            { spanContext = SpanContext span_id (TId 0),-              spanOperation = "unknown",-              spanStartedAt = 0,-              spanFinishedAt = timestamp,-              spanTags = mempty,-              spanEvents = mempty,-              spanStatus = OK,-              spanParentId = Nothing-            }+modifySpan :: HasCallStack => SpanId -> (Span -> Span) -> State -> State+modifySpan sid f st = st {spans = HM.adjust f sid (spans st)} -pushSpan :: HasCallStack => ThreadId -> T.Text -> OTel.Timestamp -> State -> State-pushSpan tid name timestamp st = st {spanStacks = new_stacks, randomGen = new_randomGen, traceMap = new_traceMap}-  where-    maybe_parent = NE.head <$> HM.lookup tid (spanStacks st)-    new_stacks = HM.alter f tid (spanStacks st)-    f Nothing = Just $ sp :| []-    f (Just sps) = Just $ cons sp sps-    (sid, new_randomGen) = R.nextWord64 (randomGen st)-    (new_traceMap, trace_id) = case (maybe_parent, HM.lookup tid (traceMap st)) of-      (Just parent, _) -> (traceMap st, spanTraceId parent)-      (_, Just trace_id') -> (traceMap st, trace_id')-      _ -> let new_trace_id = TId sid in (HM.insert tid new_trace_id (traceMap st), new_trace_id)-    sp =-      Span-        { spanContext = SpanContext (SId sid) trace_id,-          spanOperation = name,-          spanStartedAt = timestamp,-          spanFinishedAt = 0,-          spanTags = HM.singleton "tid" (IntTagValue $ fromIntegral tid),-          spanEvents = mempty,-          spanStatus = OK,-          spanParentId = spanId <$> maybe_parent-        }+-- pushSpan :: HasCallStack => ThreadId -> T.Text -> OTel.Timestamp -> Maybe Word64 -> State -> State+-- pushSpan tid name timestamp serial st = st {spanStacks = new_stacks, randomGen = new_randomGen, traceMap = new_traceMap}+--   where+--     maybe_parent = NE.head <$> HM.lookup tid (spanStacks st)+--     new_stacks = HM.alter f tid (spanStacks st)+--     f Nothing = Just $ sp :| []+--     f (Just sps) = Just $ cons sp sps+--     (sid, new_randomGen) = R.nextWord64 (randomGen st)+--     (new_traceMap, trace_id) = case (maybe_parent, HM.lookup tid (traceMap st)) of+--       (Just parent, _) -> (traceMap st, spanTraceId parent)+--       (_, Just trace_id') -> (traceMap st, trace_id')+--       _ -> let new_trace_id = TId sid in (HM.insert tid new_trace_id (traceMap st), new_trace_id)+--     sp =+--       Span+--         { spanContext = SpanContext (SId sid) trace_id,+--           spanOperation = name,+--           spanStartedAt = timestamp,+--           spanFinishedAt = 0,+--           spanTags = HM.singleton "tid" (IntTagValue $ fromIntegral tid),+--           spanEvents = mempty,+--           spanStatus = OK,+--           spanParentId = spanId <$> maybe_parent+--         } -popSpan :: HasCallStack => ThreadId -> OTel.Timestamp -> State -> (State, [Span])-popSpan tid timestamp st = (st {spanStacks = new_stacks, traceMap = new_traceMap}, [sp {spanFinishedAt = timestamp}])-  where-    sp :| new_stack = fromMaybe (error $ printf "popSpan: missing span stack for thread %d" tid) $ HM.lookup tid (spanStacks st)-    (new_traceMap, new_stacks) = case new_stack of-      [] -> (HM.delete tid (traceMap st), HM.delete tid (spanStacks st))-      x : xs -> (traceMap st, HM.insert tid (x :| xs) (spanStacks st))+-- popSpan :: HasCallStack => ThreadId -> OTel.Timestamp -> State -> (State, [Span])+-- popSpan tid timestamp st = (st {spanStacks = new_stacks, traceMap = new_traceMap}, [sp {spanFinishedAt = timestamp}])+--   where+--     sp :| new_stack = fromMaybe (error $ printf "popSpan: missing span stack for thread %d" tid) $ HM.lookup tid (spanStacks st)+--     (new_traceMap, new_stacks) = case new_stack of+--       [] -> (HM.delete tid (traceMap st), HM.delete tid (spanStacks st))+--       x : xs -> (traceMap st, HM.insert tid (x :| xs) (spanStacks st)) -pushGCSpans :: HasCallStack => State -> OTel.Timestamp -> State-pushGCSpans st timestamp = foldr go st tids-  where-    tids = HM.keys (spanStacks st)-    go tid = pushSpan tid "gc" timestamp+-- pushGCSpans :: HasCallStack => State -> OTel.Timestamp -> State+-- pushGCSpans st timestamp = foldr go st tids+--   where+--     tids = HM.keys (spanStacks st)+--     go tid = pushSpan tid "gc" timestamp Nothing -popSpansAcrossAllThreads :: HasCallStack => OTel.Timestamp -> State -> (State, [Span])-popSpansAcrossAllThreads timestamp st = foldr go (st, []) tids-  where-    tids = HM.keys (spanStacks st)-    go tid (st', sps) =-      let (st'', sps') = popSpan tid timestamp st'-       in (st'', sps' <> sps)+-- popSpansAcrossAllThreads :: HasCallStack => OTel.Timestamp -> State -> (State, [Span])+-- popSpansAcrossAllThreads timestamp st = foldr go (st, []) tids+--   where+--     tids = HM.keys (spanStacks st)+--     go tid (st', sps) =+--       let (st'', sps') = popSpan tid timestamp st'+--        in (st'', sps' <> sps)  isTerminalThreadStatus :: ThreadStopStatus -> Bool isTerminalThreadStatus HeapOverflow = True@@ -276,3 +360,17 @@  showT :: Show a => a -> T.Text showT = T.pack . show++emitSpan :: Word64 -> SpanId -> State -> (State, Span)+emitSpan serial span_id st@S {..} =+  case (HM.lookup serial serial2sid, HM.lookup span_id spans) of+    (Just span_id', Just sp)+      | span_id == span_id' ->+        ( st+            { spans = HM.delete span_id spans,+              serial2sid = HM.delete serial serial2sid,+              thread2sid = HM.update (const $ spanParentId sp) (spanThreadId sp) thread2sid+            },+          sp+        )+    _ -> error "emitSpan invariants violated"
− src/OpenTelemetry/Tracer.hs
@@ -1,43 +0,0 @@-module OpenTelemetry.Tracer where--import OpenTelemetry.Common-import qualified Data.HashMap.Strict as HM-import Data.List.NonEmpty as NE-import Data.Hashable--data Tracer threadId-  = Tracer-      { tracerSpanStacks :: !(HM.HashMap threadId (NE.NonEmpty Span))-      }-      deriving (Eq, Show)--tracerPushSpan :: (Eq tid, Hashable tid) => Tracer tid -> tid -> Span -> Tracer tid-tracerPushSpan t@(Tracer {..}) tid sp =-  case HM.lookup tid tracerSpanStacks of-    Nothing ->-      let !stacks = HM.insert tid (sp :| []) tracerSpanStacks-      in Tracer stacks-    Just sps ->-      let !stacks = HM.insert tid (sp <| sps) tracerSpanStacks-      in t { tracerSpanStacks = stacks }--tracerPopSpan :: (Eq tid, Hashable tid) => Tracer tid -> tid -> (Maybe Span, Tracer tid)-tracerPopSpan t@(Tracer {..}) tid =-  case HM.lookup tid tracerSpanStacks of-    Nothing -> (Nothing, t)-    Just (sp :| sps) ->-      let stacks =-            case NE.nonEmpty sps of-              Nothing -> HM.delete tid tracerSpanStacks-              Just sps' -> HM.insert tid sps' tracerSpanStacks-       in (Just sp, Tracer stacks)--tracerGetCurrentActiveSpan :: (Hashable tid, Eq tid) => Tracer tid -> tid -> Maybe Span-tracerGetCurrentActiveSpan (Tracer stacks) tid =-  case HM.lookup tid stacks of-    Nothing -> Nothing-    Just (sp NE.:| _) -> Just sp--createTracer :: (Hashable tid, Eq tid) => IO (Tracer tid)-createTracer = pure $ Tracer mempty-
− unit-tests/TestCommon.hs
@@ -1,36 +0,0 @@-{-# LANGUAGE OverloadedStrings #-}--module TestCommon where--import Data.Word-import OpenTelemetry.Common-import OpenTelemetry.SpanContext-import OpenTelemetry.Tracer--mkTestSpan :: Word64 -> Word64 -> Span-mkTestSpan sid tid =-  Span-    { spanContext = SpanContext (SId sid) (TId tid),-      spanStartedAt = 0,-      spanFinishedAt = 10,-      spanTags = mempty,-      spanEvents = mempty,-      spanStatus = OK,-      spanOperation = "foo",-      spanParentId = Nothing-    }--prop_tracer_push_does_something :: Word64 -> Word64 -> Int -> Bool-prop_tracer_push_does_something sid tid threadid =-  let tracer0 = Tracer mempty-      sp0 = mkTestSpan sid tid-      tracer1 = tracerPushSpan tracer0 threadid sp0-   in tracer0 /= tracer1--prop_tracer_push_pop :: Word64 -> Word64 -> Int -> Bool-prop_tracer_push_pop sid tid threadid =-  let tracer0 = Tracer mempty-      sp0 = mkTestSpan sid tid-      tracer1 = tracerPushSpan tracer0 threadid sp0-      (Just sp1, tracer2) = tracerPopSpan tracer1 threadid-   in tracer0 == tracer2 && sp0 == sp1
+ unit-tests/TestEventlogStreaming.hs view
@@ -0,0 +1,118 @@+{-# LANGUAGE OverloadedStrings #-}++module TestEventlogStreaming where++import Data.Function+import qualified Data.HashMap.Strict as HM+import qualified Data.HashSet as HS+import Data.Hashable+import Data.List (foldl', sort)+import qualified Data.Text as T+import Data.Word+import GHC.RTS.Events+import OpenTelemetry.Common hiding (Event)+import OpenTelemetry.EventlogStreaming_Internal+import OpenTelemetry.SpanContext+import Test.QuickCheck+import Text.Printf+import TextShow++instance Arbitrary SpanId where+  arbitrary = SId <$> arbitrary++processEvents :: [Event] -> State -> (State, [Span])+processEvents events st0 = foldl' go (st0, []) events+  where+    go (st, sps) e =+      let (st', sps') = processEvent e st+       in (st', sps' <> sps)++prop_number_of_spans_in_eventlog_is_number_of_spans_exported :: [(Word64, Int)] -> Bool+prop_number_of_spans_in_eventlog_is_number_of_spans_exported spans =+  let input_events = concatMap convert spans+      convert (span_serial_number, thread_id) =+        [ Event 0 (UserMessage {msg = T.pack $ printf "ot2 begin span %d %d" span_serial_number thread_id}) (Just 0),+          Event 42 (UserMessage {msg = T.pack $ printf "ot2 end span %d" span_serial_number}) (Just 0)+        ]+      (_end_state, emitted_spans) = processEvents input_events (initialState 0 (error "randomGen seed"))+   in length emitted_spans == length spans++prop_user_specified_span_ids_are_used :: [(Word64, SpanId, Int)] -> Bool+prop_user_specified_span_ids_are_used spans =+  let input_events = concatMap convert spans+      convert (span_serial_number, SId sid, thread_id) =+        [ Event 0 (UserMessage {msg = T.pack $ printf "ot2 begin span %d %d" span_serial_number thread_id}) (Just 0),+          Event 1 (UserMessage {msg = T.pack $ printf "ot2 set spanid %d %016x" span_serial_number sid}) (Just 0),+          Event 42 (UserMessage {msg = T.pack $ printf "ot2 end span %d" span_serial_number}) (Just 0)+        ]+      (_end_state, emitted_spans) = processEvents input_events (initialState 0 (error "randomGen seed"))+   in sort (map (\(_, x, _) -> x) spans) == sort (map spanId emitted_spans)++prop_user_specified_things_are_used :: [(Word64, SpanId, Int)] -> Property+prop_user_specified_things_are_used spans =+  distinct (map (\(serial, _, _) -> serial) spans)+    ==> distinct (map (\(_, span_id, _) -> span_id) spans)+    ==> classify (length spans > 1) "multiple spans"+    $ let input_events = concatMap convert spans+          convert (span_serial_number, SId sid, thread_id) =+            [ Event 0 (UserMessage {msg = T.pack $ printf "ot2 begin span %d %d" span_serial_number thread_id}) (Just 0),+              Event 1 (UserMessage {msg = T.pack $ printf "ot2 set spanid %d %016x" span_serial_number sid}) (Just 0),+              Event 2 (UserMessage {msg = T.pack $ printf "ot2 set tag %d color %d" span_serial_number sid}) (Just 0),+              Event 3 (UserMessage {msg = T.pack $ printf "ot2 set traceid %d %016x" span_serial_number sid}) (Just 0),+              Event 4 (UserMessage {msg = T.pack $ printf "ot2 add event %d message %d" span_serial_number sid}) (Just 0),+              Event 42 (UserMessage {msg = T.pack $ printf "ot2 end span %d" span_serial_number}) (Just 0)+            ]+          (_end_state, emitted_spans) = processEvents input_events (initialState 0 (error "randomGen seed"))+          corresponding_span_was_emitted (_serial, SId sid, _thread_id) =+            emitted_spans+              & filter+                ( \sp ->+                    and+                      [ spanId sp == SId sid,+                        spanTraceId sp == TId sid,+                        HM.lookup "color" (spanTags sp) == Just (StringTagValue (showt sid)),+                        any+                          (\SpanEvent {..} -> spanEventKey == "message" && spanEventValue == showt sid)+                          (spanEvents sp)+                      ]+                )+              & length+              & (== (1 :: Int))+       in conjoin $ map corresponding_span_was_emitted spans++prop_parenting_works_when_everything_is_on_one_thread_and_nested_properly :: [Word64] -> Property+prop_parenting_works_when_everything_is_on_one_thread_and_nested_properly serials =+  distinct serials ==> length serials > 1+    ==> let input_events = prelude <> map convert_begin serials <> map convert_end (reverse serials)+            prelude = [Event 0 (CreateThread 1) (Just 0)]+            convert_begin serial =+              Event 1 (UserMessage {msg = T.pack $ printf "ot2 begin span %d %d" serial serial}) (Just 0)+            convert_end serial =+              Event 42 (UserMessage {msg = T.pack $ printf "ot2 end span %d" serial}) (Just 0)+            (_end_state, emitted_spans) = processEvents input_events (initialState 0 (error "randomGen seed"))+            check_relationship (sp, psp) = spanParentId sp === Just (spanId psp)+         in conjoin $ map check_relationship (zip (tail emitted_spans) emitted_spans)++prop_beginning_a_span_on_one_thread_and_ending_on_another_is_fine :: Word64 -> ThreadId -> ThreadId -> Int -> Int -> Property+prop_beginning_a_span_on_one_thread_and_ending_on_another_is_fine serial begin_tid end_tid begin_cap end_cap =+  let input_events =+        [ Event 0 (CreateThread begin_tid) (Just begin_cap),+          Event 1 (RunThread begin_tid) (Just begin_cap),+          Event 2 (UserMessage {msg = T.pack $ printf "ot2 begin span %d %d" serial serial}) (Just begin_cap),+          Event 3 (CreateThread end_tid) (Just end_cap),+          Event 4 (RunThread end_tid) (Just end_cap),+          Event 5 (UserMessage {msg = T.pack $ printf "ot2 end span %d" serial}) (Just end_cap)+        ]+      (end_state, emitted_spans) = processEvents input_events (initialState 0 (error "randomGen seed"))+   in conjoin+        [ length emitted_spans === 1,+          spanOperation (head emitted_spans) === showt serial,+          spanStartedAt (head emitted_spans) === 2,+          spanFinishedAt (head emitted_spans) === 5,+          spanThreadId (head emitted_spans) === begin_tid,+          True === null (spans end_state),+          True === null (serial2sid end_state)+        ]++distinct :: (Eq a, Hashable a) => [a] -> Bool+distinct things = length things == HS.size (foldMap HS.singleton things)