opentelemetry-extra 0.6.1 → 0.7.0
raw patch · 7 files changed
+177/−161 lines, 7 filesdep +jsonifierdep +tasty-benchdep −aesondep −gaugePVP ok
version bump matches the API change (PVP)
Dependencies added: jsonifier, tasty-bench
Dependencies removed: aeson, gauge
API changes (from Hackage documentation)
- OpenTelemetry.ChromeExporter: ChromeTagValue :: TagValue -> ChromeTagValue
- OpenTelemetry.ChromeExporter: data DoWeCollapseThreads
- OpenTelemetry.ChromeExporter: instance Data.Aeson.Types.ToJSON.ToJSON OpenTelemetry.ChromeExporter.ChromeBeginSpan
- OpenTelemetry.ChromeExporter: instance Data.Aeson.Types.ToJSON.ToJSON OpenTelemetry.ChromeExporter.ChromeEndSpan
- OpenTelemetry.ChromeExporter: instance Data.Aeson.Types.ToJSON.ToJSON OpenTelemetry.ChromeExporter.ChromeEvent
- OpenTelemetry.ChromeExporter: instance Data.Aeson.Types.ToJSON.ToJSON OpenTelemetry.ChromeExporter.ChromeTagValue
- OpenTelemetry.ChromeExporter: newtype ChromeTagValue
- OpenTelemetry.Common: instance Data.Aeson.Types.ToJSON.ToJSON OpenTelemetry.Common.EventVal
- OpenTelemetry.Common: instance Data.Aeson.Types.ToJSON.ToJSON OpenTelemetry.Common.TagVal
- OpenTelemetry.Common: instance Data.Aeson.Types.ToJSON.ToJSONKey OpenTelemetry.Common.TagName
- OpenTelemetry.EventlogStreaming_Internal: [threadMap] :: State -> IntMap ThreadId
- OpenTelemetry.ZipkinExporter: instance Data.Aeson.Types.ToJSON.ToJSON OpenTelemetry.ZipkinExporter.ZipkinSpan
+ OpenTelemetry.ChromeExporter: data ThreadPresentation
+ OpenTelemetry.ChromeExporter: jChromeBeginSpan :: Span -> Json
+ OpenTelemetry.ChromeExporter: jChromeEndSpan :: Span -> Json
+ OpenTelemetry.ChromeExporter: jChromeEvent :: ChromeEvent -> Json
+ OpenTelemetry.ChromeExporter: jTagValue :: TagValue -> Json
+ OpenTelemetry.Common: [$sel:spanDisplayThreadId:Span] :: Span -> Word32
+ OpenTelemetry.EventlogStreaming_Internal: [cap2thread] :: State -> IntMap ThreadId
+ OpenTelemetry.EventlogStreaming_Internal: [gcRequestedAt] :: State -> !Timestamp
+ OpenTelemetry.EventlogStreaming_Internal: [nextFreeDisplayThread] :: State -> ThreadId
+ OpenTelemetry.EventlogStreaming_Internal: [thread2displayThread] :: State -> HashMap ThreadId ThreadId
+ OpenTelemetry.EventlogStreaming_Internal: inventDisplayTid :: ThreadId -> State -> (State, ThreadId)
+ OpenTelemetry.ZipkinExporter: jSpan :: ZipkinConfig -> Span -> Json
- OpenTelemetry.ChromeExporter: CollapseThreads :: DoWeCollapseThreads
+ OpenTelemetry.ChromeExporter: CollapseThreads :: ThreadPresentation
- OpenTelemetry.ChromeExporter: SplitThreads :: DoWeCollapseThreads
+ OpenTelemetry.ChromeExporter: SplitThreads :: ThreadPresentation
- OpenTelemetry.ChromeExporter: createChromeExporter' :: FilePath -> DoWeCollapseThreads -> IO (Exporter Span, Exporter Metric)
+ OpenTelemetry.ChromeExporter: createChromeExporter' :: FilePath -> ThreadPresentation -> IO (Exporter Span, Exporter Metric)
- OpenTelemetry.ChromeExporter: eventlogToChrome :: FilePath -> FilePath -> DoWeCollapseThreads -> IO ()
+ OpenTelemetry.ChromeExporter: eventlogToChrome :: FilePath -> FilePath -> ThreadPresentation -> IO ()
- OpenTelemetry.Common: Span :: {-# UNPACK #-} !SpanContext -> Text -> Word32 -> !Timestamp -> !Timestamp -> !HashMap TagName TagValue -> [SpanEvent] -> !SpanStatus -> Maybe SpanId -> !Word64 -> Span
+ OpenTelemetry.Common: Span :: {-# UNPACK #-} !SpanContext -> Text -> Word32 -> Word32 -> !Timestamp -> !Timestamp -> !HashMap TagName TagValue -> [SpanEvent] -> !SpanStatus -> Maybe SpanId -> !Word64 -> Span
- OpenTelemetry.EventlogStreaming_Internal: S :: !Timestamp -> IntMap ThreadId -> HashMap SpanId Span -> HashMap InstrumentId CaptureInstrument -> HashMap ThreadId TraceId -> HashMap Word64 SpanId -> HashMap ThreadId SpanId -> !Timestamp -> !Int -> !Int -> !Int -> !Int -> SMGen -> State
+ OpenTelemetry.EventlogStreaming_Internal: S :: !Timestamp -> IntMap ThreadId -> HashMap SpanId Span -> HashMap InstrumentId CaptureInstrument -> HashMap ThreadId TraceId -> HashMap Word64 SpanId -> HashMap ThreadId SpanId -> HashMap ThreadId ThreadId -> ThreadId -> !Timestamp -> !Timestamp -> !Int -> !Int -> !Int -> !Int -> SMGen -> State
- OpenTelemetry.ZipkinExporter: tagValue2text :: TagValue -> Text
+ OpenTelemetry.ZipkinExporter: tagValue2text :: TagValue -> Json
Files
- exe/eventlog-to-tracy/Main.hs +2/−2
- exe/ot-write-benchmark/Main.hs +1/−1
- opentelemetry-extra.cabal +4/−4
- src/OpenTelemetry/ChromeExporter.hs +68/−70
- src/OpenTelemetry/Common.hs +4/−4
- src/OpenTelemetry/EventlogStreaming_Internal.hs +63/−23
- src/OpenTelemetry/ZipkinExporter.hs +35/−57
exe/eventlog-to-tracy/Main.hs view
@@ -39,12 +39,12 @@ ["-h"] -> help ["--help"] -> help ["help"] -> help- [eventlogFile] -> work eventlogFile CollapseThreads+ [eventlogFile] -> work eventlogFile SplitThreads ["--collapse-threads", eventlogFile] -> work eventlogFile CollapseThreads ["--split-threads", eventlogFile] -> work eventlogFile SplitThreads _ -> help -work :: FilePath -> DoWeCollapseThreads -> IO ()+work :: FilePath -> ThreadPresentation -> IO () work inputFile doWeCollapseThreads = do let chromeFile = inputFile ++ ".trace.json" tracyFile = inputFile ++ ".tracy"
exe/ot-write-benchmark/Main.hs view
@@ -3,7 +3,7 @@ module Main where import Data.Word-import Gauge.Main+import Test.Tasty.Bench import qualified OpenTelemetry.Eventlog as BE import OpenTelemetry.SpanContext
opentelemetry-extra.cabal view
@@ -2,7 +2,7 @@ name: opentelemetry-extra description: The OpenTelemetry Haskell Client https://opentelemetry.io category: OpenTelemetry-version: 0.6.1+version: 0.7.0 license-file: LICENSE license: Apache-2.0 author: Dmitry Ivanov@@ -50,7 +50,6 @@ build-depends: base >= 4.12 && < 5, binary >= 0.8.6.0,- aeson, async, bytestring, clock >= 0.8,@@ -63,13 +62,14 @@ http-client, http-client-tls, http-types,+ jsonifier, opentelemetry >= 0.6.1, random >= 1.1, scientific,- text-show, splitmix, stm, text,+ text-show, unordered-containers hs-source-dirs: src@@ -141,7 +141,7 @@ -eventlog build-depends: base,- gauge >= 0.2.4,+ tasty-bench >= 0.2.4, opentelemetry >= 0.6.1, executable eventlog-to-chrome
src/OpenTelemetry/ChromeExporter.hs view
@@ -3,13 +3,14 @@ module OpenTelemetry.ChromeExporter where import Control.Monad-import Data.Aeson-import qualified Data.ByteString.Lazy as LBS+import qualified Data.ByteString as BS+import Data.Coerce import Data.Function import Data.HashMap.Strict as HM import Data.List (sortOn) import qualified Data.Text.Encoding as TE import Data.Word+import qualified Jsonifier as J import OpenTelemetry.Common import OpenTelemetry.EventlogStreaming_Internal import System.IO@@ -18,80 +19,77 @@ newtype ChromeEndSpan = ChromeEnd Span -newtype ChromeTagValue = ChromeTagValue TagValue- data ChromeEvent = ChromeEvent Word32 SpanEvent -instance ToJSON ChromeTagValue where- toJSON (ChromeTagValue (StringTagValue (TagVal i))) = Data.Aeson.String i- toJSON (ChromeTagValue (IntTagValue i)) = Data.Aeson.Number $ fromIntegral i- toJSON (ChromeTagValue (BoolTagValue b)) = Data.Aeson.Bool b- toJSON (ChromeTagValue (DoubleTagValue d)) = Data.Aeson.Number $ realToFrac d+jTagValue :: TagValue -> J.Json+jTagValue (StringTagValue (TagVal i)) = J.textString i+jTagValue (IntTagValue i) = J.intNumber i+jTagValue (BoolTagValue b) = J.bool b+jTagValue (DoubleTagValue d) = J.doubleNumber d -instance ToJSON ChromeEvent where- toJSON (ChromeEvent threadId SpanEvent {..}) =- object- [ "ph" .= ("i" :: String),- "name" .= spanEventValue,- "pid" .= (1 :: Int),- "tid" .= threadId,- "ts" .= (div spanEventTimestamp 1000)- ]+jChromeEvent (ChromeEvent threadId SpanEvent {..}) =+ J.object+ [ ("ph", J.textString "i"),+ ("name", J.textString $ coerce spanEventValue),+ ("pid", J.intNumber 1),+ ("tid", J.wordNumber $ fromIntegral threadId),+ ("ts", J.intNumber . fromIntegral $ div spanEventTimestamp 1000)+ ] -instance ToJSON ChromeBeginSpan where- toJSON (ChromeBegin Span {..}) =- object- [ "ph" .= ("B" :: String),- "name" .= spanOperation,- "pid" .= (1 :: Int),- "tid" .= spanThreadId,- "ts" .= (div spanStartedAt 1000),- "args"- .= fmap- ChromeTagValue- ( spanTags- & HM.insert "gc_us" (IntTagValue . fromIntegral $ spanNanosecondsSpentInGC `div` 1000)- & ( if spanNanosecondsSpentInGC == 0- then id- else HM.insert "gc_fraction" (DoubleTagValue (fromIntegral spanNanosecondsSpentInGC / fromIntegral (spanFinishedAt - spanStartedAt)))- )- )- ]+jChromeBeginSpan Span {..} =+ J.object+ [ ("ph", J.textString "B"),+ ("name", J.textString spanOperation),+ ("pid", J.intNumber 1),+ ("tid", J.intNumber $ fromIntegral spanDisplayThreadId),+ ("ts", J.wordNumber . fromIntegral $ div spanStartedAt 1000),+ ( "args",+ J.object+ ( spanTags+ & HM.insert "gc_us" (IntTagValue . fromIntegral $ spanNanosecondsSpentInGC `div` 1000)+ & ( if spanNanosecondsSpentInGC == 0+ then id+ else HM.insert "gc_fraction" (DoubleTagValue (fromIntegral spanNanosecondsSpentInGC / fromIntegral (spanFinishedAt - spanStartedAt)))+ )+ & HM.toList+ & fmap (\(TagName n, v) -> (n, jTagValue v))+ )+ )+ ] -instance ToJSON ChromeEndSpan where- toJSON (ChromeEnd Span {..}) =- object- [ "ph" .= ("E" :: String),- "name" .= spanOperation,- "pid" .= (1 :: Int),- "tid" .= spanThreadId,- "ts" .= (div spanFinishedAt 1000)- ]+jChromeEndSpan Span {..} =+ J.object+ [ ("ph", J.textString "E"),+ ("name", J.textString spanOperation),+ ("pid", J.intNumber 1),+ ("tid", J.intNumber $ fromIntegral spanDisplayThreadId),+ ("ts", J.wordNumber . fromIntegral $ div spanFinishedAt 1000)+ ] createChromeExporter :: FilePath -> IO (Exporter Span, Exporter Metric) createChromeExporter path = createChromeExporter' path SplitThreads -createChromeExporter' :: FilePath -> DoWeCollapseThreads -> IO (Exporter Span, Exporter Metric)-createChromeExporter' path doWeCollapseThreads = do+createChromeExporter' :: FilePath -> ThreadPresentation -> IO (Exporter Span, Exporter Metric)+createChromeExporter' path threadPresentation = do f <- openFile path WriteMode hPutStrLn f "[ "- let modifyThreadId = case doWeCollapseThreads of- CollapseThreads -> const 1- SplitThreads -> id+ let modifyThreadId = case threadPresentation of+ CollapseThreads -> pure . const 1+ SplitThreads -> pure span_exporter = Exporter ( \sps -> do mapM_- ( \sp -> do- let sp' = sp {spanThreadId = modifyThreadId (spanThreadId sp)}- let Span {spanThreadId, spanEvents} = sp'- LBS.hPutStr f $ encode $ ChromeBegin sp'- LBS.hPutStr f ",\n"+ ( \sp@(Span {spanEvents}) -> do+ tid' <- modifyThreadId (spanDisplayThreadId sp)+ let sp' = sp {spanDisplayThreadId = tid'}+ BS.hPutStr f $ J.toByteString $ jChromeBeginSpan sp'+ BS.hPutStr f ",\n" forM_ (sortOn spanEventTimestamp spanEvents) $ \ev -> do- LBS.hPutStr f $ encode $ ChromeEvent (modifyThreadId spanThreadId) ev- LBS.hPutStr f ",\n"- LBS.hPutStr f $ encode $ ChromeEnd sp'- LBS.hPutStr f ",\n"+ BS.hPutStr f $ J.toByteString $ jChromeEvent $ ChromeEvent tid' ev+ BS.hPutStr f ",\n"+ BS.hPutStr f $ J.toByteString $ jChromeEndSpan sp'+ BS.hPutStr f ",\n" ) sps pure ExportSuccess@@ -107,23 +105,23 @@ ( \metrics -> do -- forM_ metrics $ \(AggregatedMetric (SomeInstrument (TE.decodeUtf8 . instrumentName -> name)) (MetricDatapoint ts value)) -> do forM_ metrics $ \(AggregatedMetric (CaptureInstrument _ (TE.decodeUtf8 -> name)) (MetricDatapoint ts value)) -> do- LBS.hPutStr f $- encode $- object- [ "ph" .= ("C" :: String),- "name" .= name,- "ts" .= (div ts 1000),- "args" .= object [name .= Number (fromIntegral value)]+ BS.hPutStr f $+ J.toByteString $+ J.object+ [ ("ph", J.textString "C"),+ ("name", J.textString name),+ ("ts", J.wordNumber $ fromIntegral $ div ts 1000),+ ("args", J.object [(name, J.intNumber value)]) ]- LBS.hPutStr f ",\n"+ BS.hPutStr f ",\n" pure ExportSuccess ) (pure ()) pure (span_exporter, metric_exporter) -data DoWeCollapseThreads = CollapseThreads | SplitThreads+data ThreadPresentation = CollapseThreads | SplitThreads -eventlogToChrome :: FilePath -> FilePath -> DoWeCollapseThreads -> IO ()+eventlogToChrome :: FilePath -> FilePath -> ThreadPresentation -> IO () eventlogToChrome eventlogFile chromeFile doWeCollapseThreads = do (span_exporter, metric_exporter) <- createChromeExporter' chromeFile doWeCollapseThreads exportEventlog span_exporter metric_exporter eventlogFile
src/OpenTelemetry/Common.hs view
@@ -7,7 +7,6 @@ module OpenTelemetry.Common where import Control.Monad-import Data.Aeson import qualified Data.ByteString as BS import qualified Data.HashMap.Strict as HM import Data.Hashable@@ -25,13 +24,13 @@ newtype SpanName = SpanName T.Text deriving (Show, Eq, Generic) -newtype TagName = TagName T.Text deriving (Show, Eq, Generic, ToJSONKey, Hashable)+newtype TagName = TagName T.Text deriving (Show, Eq, Generic, Hashable) -newtype TagVal = TagVal T.Text deriving (Show, Eq, Generic, ToJSON)+newtype TagVal = TagVal T.Text deriving (Show, Eq, Generic) newtype EventName = EventName T.Text deriving (Show, Eq, Generic) -newtype EventVal = EventVal T.Text deriving (Show, Eq, Generic, ToJSON)+newtype EventVal = EventVal T.Text deriving (Show, Eq, Generic) instance IsString TagName where fromString = TagName . T.pack@@ -65,6 +64,7 @@ { spanContext :: {-# UNPACK #-} !SpanContext, spanOperation :: T.Text, spanThreadId :: Word32,+ spanDisplayThreadId :: Word32, spanStartedAt :: !Timestamp, spanFinishedAt :: !Timestamp, spanTags :: !(HM.HashMap TagName TagValue),
src/OpenTelemetry/EventlogStreaming_Internal.hs view
@@ -35,12 +35,15 @@ data State = S { originTimestamp :: !Timestamp,- threadMap :: IM.IntMap ThreadId,+ cap2thread :: IM.IntMap ThreadId, spans :: HM.HashMap SpanId Span, instrumentMap :: HM.HashMap InstrumentId CaptureInstrument, traceMap :: HM.HashMap ThreadId TraceId, serial2sid :: HM.HashMap Word64 SpanId, thread2sid :: HM.HashMap ThreadId SpanId,+ thread2displayThread :: HM.HashMap ThreadId ThreadId, -- https://github.com/ethercrow/opentelemetry-haskell/issues/40+ nextFreeDisplayThread :: ThreadId,+ gcRequestedAt :: !Timestamp, gcStartedAt :: !Timestamp, gcGeneration :: !Int, counterEventsProcessed :: !Int,@@ -51,7 +54,7 @@ deriving (Show) initialState :: Word64 -> R.SMGen -> State-initialState timestamp = S timestamp mempty mempty mempty mempty mempty mempty 0 0 0 0 0+initialState timestamp = S timestamp mempty mempty mempty mempty mempty mempty mempty 1 0 0 0 0 0 0 data EventSource = EventLogHandle Handle WatDoOnEOF@@ -141,9 +144,9 @@ parseOpenTelemetry _ = Nothing processEvent :: Event -> State -> (State, [Span], [Metric])-processEvent (Event ts ev m_cap) st@(S {..}) =+processEvent (Event ts ev m_cap) st@S {..} = let now = originTimestamp + ts- m_thread_id = m_cap >>= flip IM.lookup threadMap+ m_thread_id = m_cap >>= flip IM.lookup cap2thread m_trace_id = m_thread_id >>= flip HM.lookup traceMap in case (ev, m_cap, m_thread_id) of (WallClockTime {sec, nsec}, _, _) ->@@ -157,39 +160,64 @@ [Metric threadsI [MetricDatapoint now 1]] ) (RunThread tid, Just cap, _) ->- (st {threadMap = IM.insert cap tid threadMap}, [], [])+ (st {cap2thread = IM.insert cap tid cap2thread}, [], []) (StopThread tid tstatus, Just cap, _) | isTerminalThreadStatus tstatus ->- ( st- { threadMap = IM.delete cap threadMap,- traceMap = HM.delete tid traceMap- },- [],- [Metric threadsI [MetricDatapoint now (-1)]]- )+ let (t2dt, nfdt) = case HM.lookup tid thread2displayThread of+ Nothing -> (thread2displayThread, nextFreeDisplayThread)+ Just _ -> (HM.delete tid thread2displayThread, nextFreeDisplayThread - 1)+ in ( st+ { cap2thread = IM.delete cap cap2thread,+ traceMap = HM.delete tid traceMap,+ thread2displayThread = t2dt,+ nextFreeDisplayThread = nfdt+ },+ [],+ [Metric threadsI [MetricDatapoint now (-1)]]+ )+ (RequestSeqGC, _, _) ->+ (st {gcRequestedAt = now}, [], [])+ (RequestParGC, _, _) ->+ (st {gcRequestedAt = now}, [], []) (StartGC, _, _) -> (st {gcStartedAt = now}, [], []) (HeapLive {liveBytes}, _, _) -> (st, [], [Metric heapLiveBytesI [MetricDatapoint now $ fromIntegral liveBytes]])- (HeapAllocated {allocBytes}, (Just cap), _) ->+ (HeapAllocated {allocBytes}, Just cap, _) -> (st, [], [Metric (heapAllocBytesI cap) [MetricDatapoint now $ fromIntegral allocBytes]]) (EndGC, _, _) ->- let (span_id, randomGen') = R.nextWord64 randomGen- sp =+ let (gc_span_id, randomGen') = R.nextWord64 randomGen+ (gc_sync_span_id, randomGen'') = R.nextWord64 randomGen'+ sp_gc = Span { spanOperation = "gc",- spanContext = SpanContext (SId span_id) (TId span_id),+ spanContext = SpanContext (SId gc_span_id) (TId gc_span_id), spanStartedAt = gcStartedAt, spanFinishedAt = now, spanThreadId = maxBound,+ spanDisplayThreadId = maxBound, spanTags = mempty,- spanEvents = mempty,+ spanEvents = [], spanParentId = Nothing, spanStatus = OK, spanNanosecondsSpentInGC = now - gcStartedAt }- spans' = fmap (\live_span -> live_span {spanNanosecondsSpentInGC = (now - gcStartedAt) + spanNanosecondsSpentInGC live_span}) spans- st' = st {randomGen = randomGen', spans = spans'}- in (st', [sp], [Metric gcTimeI [MetricDatapoint now (fromIntegral $ now - gcStartedAt)]])+ sp_sync =+ Span+ { spanOperation = "gc_sync",+ spanContext = SpanContext (SId gc_sync_span_id) (TId gc_sync_span_id),+ spanStartedAt = gcRequestedAt,+ spanFinishedAt = gcStartedAt,+ spanThreadId = maxBound,+ spanDisplayThreadId = maxBound,+ spanTags = mempty,+ spanEvents = [],+ spanParentId = Nothing,+ spanStatus = OK,+ spanNanosecondsSpentInGC = gcStartedAt - gcRequestedAt+ }+ spans' = fmap (\live_span -> live_span {spanNanosecondsSpentInGC = (now - gcRequestedAt) + spanNanosecondsSpentInGC live_span}) spans+ st' = st {randomGen = randomGen'', spans = spans'}+ in (st', [sp_sync, sp_gc], [Metric gcTimeI [MetricDatapoint now (fromIntegral $ now - gcStartedAt)]]) (parseOpenTelemetry -> Just ev', _, fromMaybe 1 -> tid) -> handleOpenTelemetryEventlogEvent ev' st (tid, now, m_trace_id) _ -> (st, [], [])@@ -265,12 +293,14 @@ case HM.lookup serial $ serial2sid st of Nothing -> let (st', span_id) = inventSpanId serial st+ (st'', display_tid) = inventDisplayTid tid st' parent = HM.lookup tid (thread2sid st) sp = Span { spanContext = SpanContext span_id (fromMaybe (TId 42) m_trace_id), spanOperation = "", spanThreadId = tid,+ spanDisplayThreadId = display_tid, spanStartedAt = 0, spanFinishedAt = now, spanTags = mempty,@@ -279,7 +309,7 @@ spanNanosecondsSpentInGC = 0, spanParentId = parent }- in (createSpan span_id sp st', [], [])+ in (createSpan span_id sp st'', [], []) Just span_id -> let (st', sp) = emitSpan serial span_id st in (st', [sp {spanFinishedAt = now}], [])@@ -288,11 +318,13 @@ Nothing -> let (st', span_id) = inventSpanId serial st parent = HM.lookup tid (thread2sid st)+ (st'', display_tid) = inventDisplayTid tid st' sp = Span { spanContext = SpanContext span_id (fromMaybe (TId 42) m_trace_id), spanOperation = operation, spanThreadId = tid,+ spanDisplayThreadId = display_tid, spanStartedAt = now, spanFinishedAt = 0, spanTags = mempty,@@ -301,10 +333,10 @@ spanNanosecondsSpentInGC = 0, spanParentId = parent }- in (createSpan span_id sp st', [], [])+ in (createSpan span_id sp st'', [], []) Just span_id -> let (st', sp) = emitSpan serial span_id st- in (st', [sp {spanOperation = operation, spanStartedAt = now, spanThreadId = tid}], [])+ in (st', [sp {spanOperation = operation, spanStartedAt = now}], []) DeclareInstrumentEv iType iId iName -> (st {instrumentMap = HM.insert iId (CaptureInstrument iType iName) (instrumentMap st)}, [], []) MetricCaptureEv instrumentId val -> case HM.lookup instrumentId (instrumentMap st) of@@ -376,6 +408,14 @@ S {serial2sid, randomGen} = st (SId -> sid, randomGen') = R.nextWord64 randomGen st' = st {serial2sid = HM.insert serial sid serial2sid, randomGen = randomGen'}++inventDisplayTid :: ThreadId -> State -> (State, ThreadId)+inventDisplayTid tid st@(S {thread2displayThread, nextFreeDisplayThread}) =+ case HM.lookup tid thread2displayThread of+ Nothing ->+ let new_dtid = nextFreeDisplayThread+ in (st {thread2displayThread = HM.insert tid new_dtid thread2displayThread, nextFreeDisplayThread = new_dtid + 1}, new_dtid)+ Just dtid -> (st, dtid) parseText :: [T.Text] -> Maybe OpenTelemetryEventlogEvent parseText =
src/OpenTelemetry/ZipkinExporter.hs view
@@ -8,10 +8,11 @@ import Control.Concurrent.Async import Control.Concurrent.STM import Control.Monad.IO.Class-import Data.Aeson+import Data.Coerce import qualified Data.HashMap.Strict as HM import Data.Scientific import qualified Data.Text as T+import qualified Jsonifier as J import Network.HTTP.Client import Network.HTTP.Client.TLS import Network.HTTP.Types@@ -26,66 +27,43 @@ zsSpan :: Span } -tagValue2text :: TagValue -> T.Text-tagValue2text tv = case tv of+tagValue2text :: TagValue -> J.Json+tagValue2text tv = J.textString $ case tv of (StringTagValue (TagVal s)) -> s (BoolTagValue b) -> if b then "true" else "false" (IntTagValue i) -> T.pack $ show i (DoubleTagValue d) -> T.pack $ show (fromFloatDigits d) -instance ToJSON ZipkinSpan where- -- FIXME(divanov): deduplicate- toJSON (ZipkinSpan ZipkinConfig {..} s@(Span {..})) =- let TId tid = spanTraceId s- SId sid = spanId s- ts = spanStartedAt `div` 1000- duration = (spanFinishedAt - spanStartedAt) `div` 1000- in object $- [ "name" .= spanOperation,- "traceId" .= T.pack (printf "%016x" tid),- "id" .= T.pack (printf "%016x" sid),- "timestamp" .= ts,- "duration" .= duration,- "localEndpoint" .= object ["serviceName" .= zServiceName],- "tags"- .= object- ( [k .= v | (k, v) <- zGlobalTags]- <> [k .= tagValue2text v | ((TagName k), v) <- HM.toList spanTags]- ),- "annotations"- .= [ object ["timestamp" .= (t `div` 1000), "value" .= v]- | SpanEvent t _ v <- spanEvents- ]- ]- <> (maybe [] (\(SId psid) -> ["parentId" .= psid]) spanParentId)-- toEncoding (ZipkinSpan ZipkinConfig {..} s@(Span {..})) =- let TId tid = spanTraceId s- SId sid = spanId s- ts = spanStartedAt `div` 1000- duration = (spanFinishedAt - spanStartedAt) `div` 1000- in pairs- ( "name" .= spanOperation- <> "traceId" .= T.pack (printf "%016x" tid)- <> "id" .= T.pack (printf "%016x" sid)- <> "timestamp" .= ts- <> "duration" .= duration- <> "localEndpoint" .= object ["serviceName" .= zServiceName]- <> "tags"- .= object- ( [k .= v | (k, v) <- zGlobalTags]- <> [k .= tagValue2text v | ((TagName k), v) <- HM.toList spanTags]- )- <> ( maybe- mempty- (\(SId psid) -> "parentId" .= T.pack (printf "%016x" psid))- spanParentId- )- <> "annotations"- .= [ object ["timestamp" .= (t `div` 1000), "value" .= v]- | SpanEvent t _ v <- spanEvents- ]+jSpan :: ZipkinConfig -> Span -> J.Json+jSpan ZipkinConfig {..} s@(Span {..}) =+ let TId tid = spanTraceId s+ SId sid = spanId s+ ts = spanStartedAt `div` 1000+ duration = (spanFinishedAt - spanStartedAt) `div` 1000+ in J.object $+ [ ("name", J.textString spanOperation),+ ("traceId", J.textString $ T.pack (printf "%016x" tid)),+ ("id", J.textString $ T.pack (printf "%016x" sid)),+ ("timestamp", J.wordNumber $ fromIntegral ts),+ ("duration", J.wordNumber $ fromIntegral duration),+ ("localEndpoint", J.object [("serviceName", J.textString zServiceName)]),+ ( "tags",+ J.object+ ( (fmap J.textString <$> zGlobalTags)+ <> [(k, tagValue2text v) | ((TagName k), v) <- HM.toList spanTags]+ )+ ),+ ( "annotations",+ J.array+ [ J.object+ [ ("timestamp", J.wordNumber $ fromIntegral (t `div` 1000)),+ ("value", J.textString $ coerce v)+ ]+ | SpanEvent t _ v <- spanEvents+ ] )+ ]+ <> (maybe [] (\(SId psid) -> [("parentId", J.wordNumber $ fromIntegral psid)]) spanParentId) data ZipkinConfig = ZipkinConfig { zEndpoint :: String,@@ -160,11 +138,11 @@ reportSpans :: String -> Manager -> ZipkinConfig -> [Span] -> IO () reportSpans endpoint httpManager cfg sps = do dd_ "reportSpans" sps- let body = encode (map (ZipkinSpan cfg) sps)+ let body = J.toByteString $ J.array (map (jSpan cfg) sps) request = (parseRequest_ endpoint) { method = "POST",- requestBody = RequestBodyLBS body,+ requestBody = RequestBodyBS body, requestHeaders = [("Content-Type", "application/json")] } resp <- httpLbs request httpManager