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 +1/−1
- exe/eventlog-to-zipkin/Main.hs +19/−2
- opentelemetry-extra.cabal +18/−11
- src/OpenTelemetry/ChromeExporter.hs +10/−19
- src/OpenTelemetry/Common.hs +23/−26
- src/OpenTelemetry/EventlogStreaming_Internal.hs +240/−142
- src/OpenTelemetry/Tracer.hs +0/−43
- unit-tests/TestCommon.hs +0/−36
- unit-tests/TestEventlogStreaming.hs +118/−0
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)