< Summary

Information
Class: Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Stores.RelationalStructuredLogStore
Assembly: Elsa.Diagnostics.StructuredLogs.Persistence.Relational
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Diagnostics.StructuredLogs.Persistence.Relational/Stores/RelationalStructuredLogStore.cs
Line coverage
6%
Covered lines: 6
Uncovered lines: 83
Coverable lines: 89
Total lines: 141
Line coverage: 6.7%
Branch coverage
0%
Covered branches: 0
Total branches: 44
Branch coverage: 0%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
WriteAsync()0%620%
WriteManyAsync()0%110100%
QueryAsync()0%110100%
ListSourcesAsync()0%156120%
InsertAsync()0%620%
CreateCommand(...)0%7280%
CreateParameters(...)100%210%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Diagnostics.StructuredLogs.Persistence.Relational/Stores/RelationalStructuredLogStore.cs

#LineLine coverage
 1using System.Data;
 2using System.Data.Common;
 3using Elsa.Diagnostics.StructuredLogs.Contracts;
 4using Elsa.Diagnostics.StructuredLogs.Models;
 5using Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Contracts;
 6using Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Models;
 7using Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Services;
 8
 9namespace Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Stores;
 10
 111public class RelationalStructuredLogStore(
 112    IRelationalStructuredLogConnectionFactory connectionFactory,
 113    IRelationalStructuredLogDialect dialect,
 114    RelationalStructuredLogSqlBuilder sqlBuilder,
 115    RelationalStructuredLogMapper mapper,
 116    IStructuredLogSourceRegistry sourceRegistry) : IStructuredLogStore
 17{
 18    public async ValueTask WriteAsync(StructuredLogEvent logEvent, CancellationToken cancellationToken = default)
 19    {
 020        await using var connection = await connectionFactory.OpenConnectionAsync(cancellationToken);
 021        await InsertAsync(connection, mapper.Map(logEvent), null, cancellationToken);
 022        sourceRegistry.MarkSeen(logEvent.SourceId, logEvent.ReceivedAt);
 023    }
 24
 25    public async ValueTask WriteManyAsync(IReadOnlyCollection<StructuredLogEvent> logEvents, CancellationToken cancellat
 26    {
 027        if (logEvents.Count == 0)
 028            return;
 29
 030        await using var connection = await connectionFactory.OpenConnectionAsync(cancellationToken);
 031        await using var transaction = await connection.BeginTransactionAsync(cancellationToken);
 32
 033        foreach (var logEvent in logEvents)
 034            await InsertAsync(connection, mapper.Map(logEvent), transaction, cancellationToken);
 35
 036        await transaction.CommitAsync(cancellationToken);
 37
 038        foreach (var logEvent in logEvents)
 039            sourceRegistry.MarkSeen(logEvent.SourceId, logEvent.ReceivedAt);
 040    }
 41
 42    public async ValueTask<RecentStructuredLogsResult> QueryAsync(StructuredLogFilter filter, CancellationToken cancella
 43    {
 044        var query = sqlBuilder.BuildQuery(filter);
 045        await using var connection = await connectionFactory.OpenConnectionAsync(cancellationToken);
 046        await using var command = CreateCommand(connection, query.Sql, query.Parameters);
 047        await using var reader = await command.ExecuteReaderAsync(cancellationToken);
 048        var items = new List<StructuredLogEvent>();
 49
 050        while (await reader.ReadAsync(cancellationToken))
 051            items.Add(mapper.Map(reader));
 52
 053        items.Reverse();
 054        return new(items, 0);
 055    }
 56
 57    public async ValueTask<IReadOnlyCollection<StructuredLogSource>> ListSourcesAsync(CancellationToken cancellationToke
 58    {
 059        var sources = sourceRegistry.List().ToDictionary(x => x.Id, StringComparer.Ordinal);
 60
 061        await using var connection = await connectionFactory.OpenConnectionAsync(cancellationToken);
 062        await using var command = CreateCommand(connection, sqlBuilder.BuildListSources());
 063        await using var reader = await command.ExecuteReaderAsync(cancellationToken);
 64
 065        while (await reader.ReadAsync(cancellationToken))
 66        {
 067            var sourceId = reader.GetString(reader.GetOrdinal("SourceId"));
 068            if (sources.ContainsKey(sourceId))
 69                continue;
 70
 071            var lastSeen = RelationalStructuredLogMapper.ParseTimestamp(reader.GetString(reader.GetOrdinal("LastSeen")))
 072            sources[sourceId] = new()
 073            {
 074                Id = sourceId,
 075                DisplayName = sourceId,
 076                MachineName = sourceId,
 077                ProcessId = 0,
 078                LastSeen = lastSeen,
 079                Status = StructuredLogSourceStatus.Connected
 080            };
 81        }
 82
 083        return sources.Values
 084            .OrderBy(x => x.DisplayName, StringComparer.OrdinalIgnoreCase)
 085            .ToList();
 086    }
 87
 88    private async ValueTask InsertAsync(DbConnection connection, RelationalStructuredLogRecord record, DbTransaction? tr
 89    {
 090        await using var command = CreateCommand(connection, sqlBuilder.BuildInsert(), CreateParameters(record));
 091        command.Transaction = transaction;
 092        await command.ExecuteNonQueryAsync(cancellationToken);
 093    }
 94
 95    private DbCommand CreateCommand(DbConnection connection, string sql, IReadOnlyDictionary<string, object?>? parameter
 96    {
 097        var command = connection.CreateCommand();
 098        command.CommandText = sql;
 099        command.CommandType = CommandType.Text;
 100
 0101        if (parameters == null)
 0102            return command;
 103
 0104        foreach (var (name, value) in parameters)
 105        {
 0106            var parameter = command.CreateParameter();
 0107            parameter.ParameterName = name.StartsWith(dialect.ParameterPrefix, StringComparison.Ordinal) ? name : $"{dia
 0108            parameter.Value = value ?? DBNull.Value;
 0109            command.Parameters.Add(parameter);
 110        }
 111
 0112        return command;
 113    }
 114
 115    private static IReadOnlyDictionary<string, object?> CreateParameters(RelationalStructuredLogRecord record)
 116    {
 0117        return new Dictionary<string, object?>
 0118        {
 0119            ["Id"] = record.Id,
 0120            ["Sequence"] = record.Sequence,
 0121            ["Timestamp"] = record.Timestamp,
 0122            ["ReceivedAt"] = record.ReceivedAt,
 0123            ["Level"] = (int)record.Level,
 0124            ["Category"] = record.Category,
 0125            ["EventId"] = record.EventId,
 0126            ["EventName"] = record.EventName,
 0127            ["Message"] = record.Message,
 0128            ["MessageTemplate"] = record.MessageTemplate,
 0129            ["ExceptionJson"] = record.ExceptionJson,
 0130            ["ScopesJson"] = record.ScopesJson,
 0131            ["PropertiesJson"] = record.PropertiesJson,
 0132            ["TraceId"] = record.TraceId,
 0133            ["SpanId"] = record.SpanId,
 0134            ["CorrelationId"] = record.CorrelationId,
 0135            ["TenantId"] = record.TenantId,
 0136            ["WorkflowDefinitionId"] = record.WorkflowDefinitionId,
 0137            ["WorkflowInstanceId"] = record.WorkflowInstanceId,
 0138            ["SourceId"] = record.SourceId
 0139        };
 140    }
 141}