Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
98 changes: 45 additions & 53 deletions src/SeqCli/Cli/Commands/IngestCommand.cs
Original file line number Diff line number Diff line change
Expand Up @@ -81,66 +81,58 @@ public IngestCommand()

protected override async Task<int> Run()
{
try
{
var enrichers = new List<IEventEnricher>();
var enrichers = new List<IEventEnricher>();

if (_level != null)
enrichers.Add(new LevelEnricher(_level));
if (_level != null)
enrichers.Add(new LevelEnricher(_level));

foreach (var (name, value) in _properties.FlatProperties)
enrichers.Add(new ScalarPropertyEnricher(name, value));
foreach (var (name, value) in _properties.FlatProperties)
enrichers.Add(new ScalarPropertyEnricher(name, value));

Func<JsonObject, bool>? filter = null;
if (_filter != null)
{
var eval = SeqSyntax.CompileExpression(_filter);
filter = evt => eval(evt).IsTrue();
}
Func<JsonObject, bool>? filter = null;
if (_filter != null)
{
var eval = SeqSyntax.CompileExpression(_filter);
filter = evt => eval(evt).IsTrue();
}

var config = RuntimeConfigurationLoader.Load(_storagePath);
var connection = SeqConnectionFactory.Connect(_connection, config);
// The API key is passed through separately because `SeqConnection` doesn't expose a batched ingestion
// mechanism and so we manually construct `HttpRequestMessage`s deeper in the stack. Nice feature gap to
// close at some point!
var (_, apiKey) = SeqConnectionFactory.GetConnectionDetails(_connection, config);
var batchSize = _batchSize.Value;
var config = RuntimeConfigurationLoader.Load(_storagePath);
var connection = SeqConnectionFactory.Connect(_connection, config);

// The API key is passed through separately because `SeqConnection` doesn't expose a batched ingestion
// mechanism and so we manually construct `HttpRequestMessage`s deeper in the stack. Nice feature gap to
// close at some point!
var (_, apiKey) = SeqConnectionFactory.GetConnectionDetails(_connection, config);
var batchSize = _batchSize.Value;

foreach (var input in _fileInputFeature.OpenInputs())
foreach (var input in _fileInputFeature.OpenInputs())
{
using (input)
{
using (input)
{
IEventReader reader = _json
? new JsonEventReader(input)
: new PlainTextEventReader(input, _pattern);

reader = new EnrichingReader(reader, enrichers);

if (_message != null)
reader = new StaticMessageTemplateReader(reader, _message);

var exit = await LogShipper.ShipEventsAsync(
connection,
apiKey,
reader,
_invalidDataHandlingFeature.InvalidDataHandling,
_sendFailureHandlingFeature.SendFailureHandling,
batchSize,
filter,
CancellationToken.None);

if (exit != 0)
return exit;
}
IEventReader reader = _json
? new JsonEventReader(input)
: new PlainTextEventReader(input, _pattern);

reader = new EnrichingReader(reader, enrichers);

if (_message != null)
reader = new StaticMessageTemplateReader(reader, _message);

var exit = await LogShipper.ShipEventsAsync(
connection,
apiKey,
reader,
_invalidDataHandlingFeature.InvalidDataHandling,
_sendFailureHandlingFeature.SendFailureHandling,
batchSize,
filter,
CancellationToken.None);

if (exit != 0)
return exit;
}

return 0;
}
catch (Exception ex)
{
Log.Error(ex, "Ingestion failed: {ErrorMessage}", ex.Message);
return 1;
}

return 0;
}
}
81 changes: 36 additions & 45 deletions src/SeqCli/Cli/Commands/Metrics/SearchCommand.cs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,6 @@
using SeqCli.Cli.Features;
using SeqCli.Config;
using SeqCli.Util;
using Serilog;

namespace SeqCli.Cli.Commands.Metrics;

Expand Down Expand Up @@ -69,56 +68,48 @@ public SearchCommand()

protected override async Task<int> Run()
{
try
{
var config = RuntimeConfigurationLoader.Load(_storagePath);
var output = _output.GetOutputFormat(config);
var connection = SeqConnectionFactory.Connect(_connection, config);
var config = RuntimeConfigurationLoader.Load(_storagePath);
var output = _output.GetOutputFormat(config);
var connection = SeqConnectionFactory.Connect(_connection, config);

string? filter = null;
if (!string.IsNullOrWhiteSpace(_filter))
filter = (await connection.Expressions.ToStrictAsync(_filter)).StrictExpression;
string? filter = null;
if (!string.IsNullOrWhiteSpace(_filter))
filter = (await connection.Expressions.ToStrictAsync(_filter)).StrictExpression;

var result = await connection.Metrics.SearchAsync(
_groups,
filter,
_count,
rangeStartUtc: _range.Start,
rangeEndUtc: _range.End,
trace: _trace);

// We convert the metric into a query result to improve formatting consistency. Room for an abstraction of
// some kind here.
var rows = new List<object?[]>();
foreach (var metric in result.Metrics)
{
var row = new List<object?>
{
metric.Name ?? metric.Accessor,
metric.Kind,
metric.Unit,
metric.Description
};

foreach (var value in metric.GroupKey)
row.Add(value);

rows.Add(row.ToArray());
}
var asRowset = new QueryResultPart
var result = await connection.Metrics.SearchAsync(
_groups,
filter,
_count,
rangeStartUtc: _range.Start,
rangeEndUtc: _range.End,
trace: _trace);

// We convert the metric into a query result to improve formatting consistency. Room for an abstraction of
// some kind here.
var rows = new List<object?[]>();
foreach (var metric in result.Metrics)
{
var row = new List<object?>
{
Columns = new[] { "Name", "Kind", "Unit", "Description" }.Concat(_groups).ToArray(),
Rows = rows.ToArray()
metric.Name ?? metric.Accessor,
metric.Kind,
metric.Unit,
metric.Description
};

output.WriteQueryResult(asRowset);

return 0;
foreach (var value in metric.GroupKey)
row.Add(value);

rows.Add(row.ToArray());
}
catch (Exception ex)
var asRowset = new QueryResultPart
{
Log.Error(ex, "Could not retrieve metrics: {ErrorMessage}", ex.Message);
return 1;
}
Columns = new[] { "Name", "Kind", "Unit", "Description" }.Concat(_groups).ToArray(),
Rows = rows.ToArray()
};

output.WriteQueryResult(asRowset);

return 0;
}
}
84 changes: 38 additions & 46 deletions src/SeqCli/Cli/Commands/SearchCommand.cs
Original file line number Diff line number Diff line change
Expand Up @@ -72,62 +72,54 @@ public SearchCommand()

protected override async Task<int> Run()
{
try
{
var config = RuntimeConfigurationLoader.Load(_storagePath);
var config = RuntimeConfigurationLoader.Load(_storagePath);

var connection = SeqConnectionFactory.Connect(_connection, config);
connection.Client.HttpClient.Timeout = TimeSpan.FromMilliseconds(_httpClientTimeout);
var connection = SeqConnectionFactory.Connect(_connection, config);
connection.Client.HttpClient.Timeout = TimeSpan.FromMilliseconds(_httpClientTimeout);

var columns = await _eventColumns.GetColumns(connection, _signal.Signal);
var output = _output.GetOutputFormat(config, TextFormatters.PlainOutputTemplate(columns));
var columns = await _eventColumns.GetColumns(connection, _signal.Signal);
var output = _output.GetOutputFormat(config, TextFormatters.PlainOutputTemplate(columns));

string? filter = null;
if (!string.IsNullOrWhiteSpace(_filter))
filter = (await connection.Expressions.ToStrictAsync(_filter)).StrictExpression;
string? filter = null;
if (!string.IsNullOrWhiteSpace(_filter))
filter = (await connection.Expressions.ToStrictAsync(_filter)).StrictExpression;

try
try
{
if (!_noWebSockets)
{
if (!_noWebSockets)
await foreach (var evt in connection.Events.EnumerateAsync(null,
_signal.Signal,
filter,
_count,
fromDateUtc: _range.Start,
toDateUtc: _range.End,
trace: _trace,
render: output.RequiresRender))
{
await foreach (var evt in connection.Events.EnumerateAsync(null,
_signal.Signal,
filter,
_count,
fromDateUtc: _range.Start,
toDateUtc: _range.End,
trace: _trace,
render: output.RequiresRender))
{
output.WriteEventEntity(evt);
}

return 0;
output.WriteEventEntity(evt);
}
}
catch (NotSupportedException nse)
{
Log.Information(nse, "WebSockets not supported; falling back to paged search");
}

await foreach (var evt in connection.Events.PagedEnumerateAsync(null,
_signal.Signal,
filter,
_count,
fromDateUtc: _range.Start,
toDateUtc: _range.End,
trace: _trace,
render: output.RequiresRender))
{
output.WriteEventEntity(evt);
}

return 0;
return 0;
}
}
catch (Exception ex)
catch (NotSupportedException nse)
{
Log.Error(ex, "Could not retrieve search result: {ErrorMessage}", ex.Message);
return 1;
Log.Information(nse, "WebSockets not supported; falling back to paged search");
}

await foreach (var evt in connection.Events.PagedEnumerateAsync(null,
_signal.Signal,
filter,
_count,
fromDateUtc: _range.Start,
toDateUtc: _range.End,
trace: _trace,
render: output.RequiresRender))
{
output.WriteEventEntity(evt);
}

return 0;
}
}
Loading