#pragma once #include "CoreMinimal.h" #include "Containers/Queue.h" #include "HAL/Runnable.h" #include "Telemetry/TelemetryEvent.h" class FRunnableThread; class FEvent; class IFileHandle; // Where events go. The emitter never knows. Three implementations on day one; a fourth adapts an analytics // provider when there is one (Docs/Spec/Telemetry.md, Q2). class SALTYCORE_API ITelemetrySink { public: virtual ~ITelemetrySink() = default; virtual void Emit(const FTelemetryEnvelope& Envelope, const FTelemetryEvent& Event) = 0; virtual void Flush() {} }; // Does nothing. The shipping default until there is somewhere to send anything. class SALTYCORE_API FNullTelemetrySink : public ITelemetrySink { public: virtual void Emit(const FTelemetryEnvelope&, const FTelemetryEvent&) override {} }; // One line of JSON to the output log under LogTelemetry. Verifies a step emitted what it claims. class SALTYCORE_API FLogTelemetrySink : public ITelemetrySink { public: virtual void Emit(const FTelemetryEnvelope& Envelope, const FTelemetryEvent& Event) override; }; // Appends one line per event to a JSON Lines file. Emit only builds the line and queues it; a worker thread // drains the queue to disk every DrainIntervalSeconds and on Flush, so the calling thread never touches the // file. Flushed on quit and from the unhandled-exception handler, so a crash loses at most one interval. class SALTYCORE_API FJsonlTelemetrySink : public ITelemetrySink, private FRunnable { public: explicit FJsonlTelemetrySink(const FString& InFilePath, double DrainIntervalSeconds = 2.0); virtual ~FJsonlTelemetrySink() override; virtual void Emit(const FTelemetryEnvelope& Envelope, const FTelemetryEvent& Event) override; virtual void Flush() override; const FString& GetFilePath() const { return FilePath; } private: // FRunnable virtual uint32 Run() override; virtual void Stop() override; void Drain(); FString FilePath; double DrainIntervalSeconds; TQueue Queue; TUniquePtr File; FEvent* WakeEvent = nullptr; FEvent* DrainedEvent = nullptr; FRunnableThread* Thread = nullptr; TAtomic bStopping{ false }; TAtomic FlushesRequested{ 0 }; TAtomic FlushesCompleted{ 0 }; FDelegateHandle OnExitHandle; FDelegateHandle OnSystemErrorHandle; };