Compare commits

...

2 Commits

Author SHA1 Message Date
Jedd Morgan 906ff9c3ff test: duck tests (#455)
.NET Build and Publish / build (push) Has been cancelled
* lets get some testing action!

* revert url

* we'll try these fixes

* we'll try this, if it doesn't work I'm deleting the test

* try this

* Add optional parent id to traces

* how about now

* and as zero

* Fix build

* fix code cov

* FML subscription tests are so flaky
2026-03-11 16:30:24 +00:00
Jedd Morgan 515d45528d feat(otel): traces for new send pipeline (#457)
* 😃 extra traces

* default
2026-03-11 14:51:58 +00:00
18 changed files with 274 additions and 47 deletions
+5 -1
View File
@@ -1,7 +1,11 @@
name: PR Test
on:
pull_request:
pull_request: {}
push:
branches:
- "main" # Need to run for codecov to compare against the BASE
jobs:
build:
-7
View File
@@ -46,13 +46,6 @@ jobs:
SEMVER: ${{ steps.set-version.outputs.SEMVER }}
FILE_VERSION: ${{ steps.set-version.outputs.FILE_VERSION }}
- name: Upload coverage reports to Codecov with GitHub Action
uses: codecov/codecov-action@v5
continue-on-error: true
with:
fail_ci_if_error: true
files: tests/**/coverage.xml
token: ${{ secrets.CODECOV_TOKEN }}
- name: NuGet login (OIDC → temp API key)
uses: NuGet/login@v1
@@ -4,5 +4,9 @@ namespace Speckle.Sdk.Logging;
public interface ISdkActivityFactory : IDisposable
{
ISdkActivity? Start(string? name = default, [CallerMemberName] string source = "");
/// <param name="name"></param>
/// <param name="source"></param>
/// <param name="parentId">Only need to set if the parent is coming from an external source (e.g.to trace between client and server)</param>
/// <returns></returns>
ISdkActivity? Start(string? name = default, [CallerMemberName] string source = "", string? parentId = null);
}
@@ -14,7 +14,8 @@ public record ModelIngestionCreateInput(
string modelId,
string projectId,
string progressMessage,
SourceDataInput sourceData
SourceDataInput sourceData,
int? maxIdleTimeoutSeconds = null
);
public record ModelIngestionUpdateInput(string ingestionId, string projectId, string progressMessage, double? progress);
@@ -4,5 +4,5 @@ public sealed class NullActivityFactory : ISdkActivityFactory
{
public void Dispose() { }
public ISdkActivity? Start(string? name = default, string source = "") => null;
public ISdkActivity? Start(string? name = default, string source = "", string? parentId = null) => null;
}
@@ -14,20 +14,20 @@ public partial interface IIngestionProgressManager : IProgress<CardProgress>;
/// An <see langword="IProgress{IngestionProgressEventArgs}"/> implementation for the entire client side Ingestion progress update reporting
/// Will throttles ingestion progress messages and reports their progress
/// </summary>
/// <remarks>
/// Normally we would pick quite a coarse updateInterval to try and spamming the server (1-5s)
/// </remarks>
[GenerateAutoInterface]
public sealed class IngestionProgressManager(
ILogger<IngestionProgressManager> logger,
IClient speckleClient,
ModelIngestion ingestion,
string projectId,
TimeSpan updateInterval,
CancellationToken cancellationToken
) : IIngestionProgressManager
{
/// <remarks>
/// Normally we would pick quite a coarse throttle window to try and avoid over pressure (1-5s)
/// </remarks>
private Task? _lastUpdate;
public Task? LastUpdate { get; private set; }
private long _lastUpdatedAt;
private readonly object _lock = new();
@@ -48,15 +48,15 @@ public sealed class IngestionProgressManager(
trimmedMessage = value.Status.TrimEnd('.');
_lastUpdate = speckleClient
LastUpdate = speckleClient
.Ingestion.UpdateProgress(
new ModelIngestionUpdateInput(ingestion.id, projectId, trimmedMessage, value.Progress),
new ModelIngestionUpdateInput(ingestion.id, ingestion.projectId, trimmedMessage, value.Progress),
cancellationToken
)
.ContinueWith(
HandleFaultedContinuation,
Continuation,
CancellationToken.None,
TaskContinuationOptions.OnlyOnFaulted | TaskContinuationOptions.ExecuteSynchronously,
TaskContinuationOptions.ExecuteSynchronously,
TaskScheduler.Default
);
}
@@ -67,7 +67,7 @@ public sealed class IngestionProgressManager(
/// <returns><see langword="true"/> if the update should be ignored, otherwise <see langword="false"/></returns>
private bool ShouldIgnoreProgressUpdate()
{
if (_lastUpdate is not null && !_lastUpdate.IsCompleted)
if (LastUpdate is not null && !LastUpdate.IsCompleted)
{
return true;
}
@@ -76,7 +76,7 @@ public sealed class IngestionProgressManager(
return msSinceLastUpdate < updateInterval;
}
private void HandleFaultedContinuation(Task updateTask)
private void Continuation(Task updateTask)
{
// The progress report failed... could be many reasons.
// For now, we're not letting this fail the Ingestion in any way
@@ -12,11 +12,10 @@ public sealed class IngestionProgressManagerFactory(ILogger<IngestionProgressMan
public IIngestionProgressManager CreateInstance(
IClient speckleClient,
ModelIngestion ingestion,
string projectId,
TimeSpan updateInterval,
CancellationToken cancellationToken
)
{
return new IngestionProgressManager(logger, speckleClient, ingestion, projectId, updateInterval, cancellationToken);
return new IngestionProgressManager(logger, speckleClient, ingestion, updateInterval, cancellationToken);
}
}
@@ -18,9 +18,9 @@ public sealed class RenderedStreamProgress(IProgress<CardProgress> progress) : I
);
}
private static readonly string[] s_suffixes = ["B", "KB", "MB", "GB", "TB", "PB", "EB", "ZB", "YB"];
private static readonly string[] s_suffixes = ["B", "KB", "MB", "GB", "TB", "PB"];
private static (string suffix, double scaleFactor) GetFileSizeRendering(long value)
internal static (string suffix, double scaleFactor) GetFileSizeRendering(long value)
{
if (value <= 0)
{
+12 -3
View File
@@ -3,13 +3,15 @@ using Microsoft.Extensions.Logging;
using Speckle.InterfaceGenerator;
using Speckle.Sdk.Dependencies;
using Speckle.Sdk.Helpers;
using Speckle.Sdk.Logging;
namespace Speckle.Sdk.Pipelines.Send;
[GenerateAutoInterface]
public sealed class DiskStoreFactory(ILogger<DiskStore> logger) : IDiskStoreFactory
public sealed class DiskStoreFactory(ILogger<DiskStore> logger, ISdkActivityFactory activityFactory) : IDiskStoreFactory
{
public DiskStore CreateInstance(CancellationToken cancellationToken) => new(logger, cancellationToken);
public DiskStore CreateInstance(CancellationToken cancellationToken) =>
new(logger, activityFactory, cancellationToken);
}
public sealed class DiskStore
@@ -17,11 +19,17 @@ public sealed class DiskStore
private readonly RepackedChannel<UploadItem> _channel;
private readonly Task<DisposableFile> _writeToDiskTask;
private readonly ILogger<DiskStore> _logger;
private readonly ISdkActivityFactory _activityFactory;
private readonly CancellationToken _cancellationToken;
internal DiskStore(ILogger<DiskStore> logger, CancellationToken cancellationToken)
internal DiskStore(
ILogger<DiskStore> logger,
ISdkActivityFactory activityFactory,
CancellationToken cancellationToken
)
{
_logger = logger;
_activityFactory = activityFactory;
_cancellationToken = cancellationToken;
_channel = new RepackedChannel<UploadItem>(1000, true, false);
@@ -33,6 +41,7 @@ public sealed class DiskStore
public async Task<DisposableFile> CompleteAsync()
{
using var a = _activityFactory.Start("Waiting for DiskStore to complete");
_channel.CompleteWriter();
return await _writeToDiskTask.ConfigureAwait(false);
}
@@ -70,11 +70,5 @@ public sealed class SendPipeline : IDisposable
await _uploader.Send(fileStreamUpload).ConfigureAwait(false);
}
public async Task<string> WaitForUploadAndServerProcessing()
{
// TODO: in some way, wait for the server to process the upload and return the actual new version id
return await Task.FromResult("todo").ConfigureAwait(false);
}
public void Dispose() => _uploader.Dispose();
}
+12 -7
View File
@@ -1,16 +1,17 @@
using System.Net.Http.Headers;
using System.Text;
using Microsoft.Extensions.Logging;
using Speckle.InterfaceGenerator;
using Speckle.Newtonsoft.Json;
using Speckle.Sdk.Credentials;
using Speckle.Sdk.Helpers;
using Speckle.Sdk.Logging;
using Speckle.Sdk.Pipelines.Progress;
namespace Speckle.Sdk.Pipelines.Send;
[GenerateAutoInterface]
public sealed class UploaderFactory(ISpeckleHttp httpClientFactory, ILogger<Uploader> logger) : IUploaderFactory
public sealed class UploaderFactory(ISpeckleHttp httpClientFactory, ISdkActivityFactory activityFactory)
: IUploaderFactory
{
public Uploader CreateInstance(
string projectId,
@@ -18,7 +19,7 @@ public sealed class UploaderFactory(ISpeckleHttp httpClientFactory, ILogger<Uplo
Account account,
IProgress<StreamProgressArgs> progress,
CancellationToken cancellationToken
) => new(projectId, ingestionId, logger, httpClientFactory, account, progress, cancellationToken);
) => new(projectId, ingestionId, activityFactory, httpClientFactory, account, progress, cancellationToken);
}
public sealed class Uploader : IDisposable
@@ -28,13 +29,13 @@ public sealed class Uploader : IDisposable
private readonly CancellationToken _cancellationToken;
private readonly HttpClient _speckleClient;
private readonly HttpClient _s3Client;
private readonly ILogger<Uploader> _logger;
private readonly ISdkActivityFactory _activity;
private readonly IProgress<StreamProgressArgs> _progress;
internal Uploader(
string projectId,
string ingestionId,
ILogger<Uploader> logger,
ISdkActivityFactory activity,
ISpeckleHttp httpClientFactory,
Account speckleAccount,
IProgress<StreamProgressArgs> progress,
@@ -43,7 +44,7 @@ public sealed class Uploader : IDisposable
{
_projectId = projectId;
_ingestionId = ingestionId;
_logger = logger;
_activity = activity;
_cancellationToken = cancellationToken;
_progress = progress;
_speckleClient = httpClientFactory.CreateHttpClient(authorizationToken: speckleAccount.token);
@@ -62,6 +63,8 @@ public sealed class Uploader : IDisposable
private async Task<PresignedUploadResponse> GetPresignedUrl()
{
using var a = _activity.Start("Get Presigned Url");
var signUri = new Uri($"projects/{_projectId}/modelingestion/{_ingestionId}/uploads/sign", UriKind.Relative);
using var signResponse = await _speckleClient.PostAsync(signUri, null, _cancellationToken).ConfigureAwait(false);
@@ -80,7 +83,7 @@ public sealed class Uploader : IDisposable
private async Task<string> UploadToS3(Stream fileStream, PresignedUploadResponse presignedUploadResponse)
{
_logger.LogInformation("Uploading file to pre-signed url");
using var a = _activity.Start("Uploading file to pre-signed url");
Stream progressStream = new ProgressStream(fileStream, _progress);
@@ -107,6 +110,8 @@ public sealed class Uploader : IDisposable
private async Task TriggerProcessing(TriggerUploadRequest request)
{
using var a = _activity.Start("Triggering Processing");
Uri processUri = new($"projects/{_projectId}/modelingestion/{_ingestionId}/uploads/process", UriKind.Relative);
string requestBody = JsonConvert.SerializeObject(request);
using var content = new StringContent(requestBody, Encoding.UTF8, "application/json");
@@ -52,7 +52,7 @@ public class ServerObjectManagerTests : MoqTest
http.Setup(x => x.CreateHttpClient(It.IsAny<HttpClientHandler>(), timeout, token)).Returns(httpClient);
var activityFactory = Create<ISdkActivityFactory>();
activityFactory.Setup(x => x.Start(null, "DownloadObjects")).Returns((ISdkActivity?)null);
activityFactory.Setup(x => x.Start(null, "DownloadObjects", null)).Returns((ISdkActivity?)null);
var serverObjectManager = new ServerObjectManager(
http.Object,
@@ -91,7 +91,7 @@ public class ServerObjectManagerTests : MoqTest
http.Setup(x => x.CreateHttpClient(It.IsAny<HttpClientHandler>(), timeout, token)).Returns(httpClient);
var activityFactory = Create<ISdkActivityFactory>();
activityFactory.Setup(x => x.Start(null, "DownloadSingleObject")).Returns((ISdkActivity?)null);
activityFactory.Setup(x => x.Start(null, "DownloadSingleObject", null)).Returns((ISdkActivity?)null);
var serverObjectManager = new ServerObjectManager(
http.Object,
@@ -132,7 +132,7 @@ public class ServerObjectManagerTests : MoqTest
http.Setup(x => x.CreateHttpClient(It.IsAny<HttpClientHandler>(), timeout, token)).Returns(httpClient);
var activityFactory = Create<ISdkActivityFactory>();
activityFactory.Setup(x => x.Start(null, "HasObjects")).Returns((ISdkActivity?)null);
activityFactory.Setup(x => x.Start(null, "HasObjects", null)).Returns((ISdkActivity?)null);
var serverObjectManager = new ServerObjectManager(
http.Object,
@@ -171,7 +171,7 @@ public class ServerObjectManagerTests : MoqTest
http.Setup(x => x.CreateHttpClient(It.IsAny<HttpClientHandler>(), timeout, token)).Returns(httpClient);
var activityFactory = Create<ISdkActivityFactory>();
activityFactory.Setup(x => x.Start(null, "UploadObjects")).Returns((ISdkActivity?)null);
activityFactory.Setup(x => x.Start(null, "UploadObjects", null)).Returns((ISdkActivity?)null);
var serverObjectManager = new ServerObjectManager(
http.Object,
@@ -13,7 +13,7 @@ public class SubscriptionResourceTests : IAsyncLifetime
#if DEBUG
private const int WAIT_PERIOD = 3000; // WSL is slow AF, so for local runs, we're being extra generous
#else
private const int WAIT_PERIOD = 400; // For CI runs, a much smaller wait time is acceptable
private const int WAIT_PERIOD = 600; // For CI runs, a much smaller wait time is acceptable
#endif
private const int TIMEOUT = WAIT_PERIOD + WAIT_PERIOD + 600;
private IClient _testUser;
@@ -0,0 +1,76 @@
using Microsoft.Extensions.DependencyInjection;
using Speckle.Sdk.Api;
using Speckle.Sdk.Api.GraphQL.Models;
using Speckle.Sdk.Common;
using Speckle.Sdk.Pipelines.Progress;
namespace Speckle.Sdk.Tests.Integration.Pipelines.Progress;
[Trait("Server", "Internal")]
public class IngestionProgressManagerTests : IAsyncLifetime
{
private IIngestionProgressManagerFactory _factory;
private IClient _client;
private Project _project;
private Model _model;
private ModelIngestion _ingestion;
public async Task InitializeAsync()
{
var serviceProvider = TestServiceSetup.GetServiceProvider();
_factory = serviceProvider.GetRequiredService<IIngestionProgressManagerFactory>();
_client = await Fixtures.SeedUserWithClient();
_project = await _client.Project.Create(new("test", null, default));
_model = await _client.Model.Create(new("test", null, _project.id));
_ingestion = await _client.Ingestion.Create(
new(_model.id, _project.id, "Testing ingestion", new("integrationTests", "0.0.0", null, null))
);
}
[Fact]
public async Task TestProgress_NoThrottle()
{
var sut = _factory.CreateInstance(_client, _ingestion, TimeSpan.Zero, CancellationToken.None);
const string FIRST_MESSAGE = "This is a test 123";
const string SECOND_MESSAGE = "This is another test 321";
// first message (should go through)
sut.Report(new CardProgress(FIRST_MESSAGE, 0.123123123d));
await sut.LastUpdate.NotNull();
var res = await _client.Ingestion.Get(_ingestion.id, _project.id, CancellationToken.None);
Assert.Equal(FIRST_MESSAGE, res.statusData.progressMessage);
// second message (should also go through)
sut.Report(new CardProgress(SECOND_MESSAGE, 0.321321321d));
await sut.LastUpdate.NotNull();
res = await _client.Ingestion.Get(_ingestion.id, _project.id, CancellationToken.None);
Assert.Equal(SECOND_MESSAGE, res.statusData.progressMessage);
}
[Fact]
public async Task TestProgress_WithThrottle()
{
var sut = _factory.CreateInstance(_client, _ingestion, TimeSpan.FromMilliseconds(500), CancellationToken.None);
const string EXPECTED_MESSAGE = "First message should go through 123";
await Task.Delay(TimeSpan.FromMilliseconds(600));
// first message (should go through)
sut.Report(new CardProgress(EXPECTED_MESSAGE, 0.123123123d));
// second message (should be dropped)
sut.Report(new CardProgress("Second message, should be dropped", 0.321321321d));
await sut.LastUpdate.NotNull();
var res = await _client.Ingestion.Get(_ingestion.id, _project.id, CancellationToken.None);
Assert.Equal(EXPECTED_MESSAGE, res.statusData.progressMessage);
}
public Task DisposeAsync()
{
_client.Dispose();
return Task.CompletedTask;
}
}
@@ -16,4 +16,7 @@
<ProjectReference Include="..\..\src\Speckle.Sdk\Speckle.Sdk.csproj" />
<ProjectReference Include="..\Speckle.Sdk.Testing\Speckle.Sdk.Testing.csproj" />
</ItemGroup>
<ItemGroup>
<Folder Include="Pipelines\Send\" />
</ItemGroup>
</Project>
@@ -0,0 +1,21 @@
using Moq;
using Speckle.Sdk.Pipelines.Progress;
namespace Speckle.Sdk.Tests.Unit.Pipelines.Progress;
public class AggregateProgressTests
{
[Fact]
public void Report_InvokesReportOnAllInnerProgresses()
{
var mock1 = new Mock<IProgress<int>>();
var mock2 = new Mock<IProgress<int>>();
const int TEST_VALUE = 42;
var target = new AggregateProgress<int>(mock1.Object, mock2.Object);
target.Report(TEST_VALUE);
mock1.Verify(x => x.Report(TEST_VALUE), Times.Once);
mock2.Verify(x => x.Report(TEST_VALUE), Times.Once);
}
}
@@ -0,0 +1,72 @@
using System.Diagnostics.CodeAnalysis;
using Moq;
using Speckle.Sdk.Pipelines.Progress;
namespace Speckle.Sdk.Tests.Unit.Pipelines.Progress;
[SuppressMessage(
"Performance",
"CA1835:Prefer the \'Memory\'-based overloads for \'ReadAsync\' and \'WriteAsync\'",
Justification = "Need to test it"
)]
public class ProgressStreamTests : IDisposable
{
private readonly Mock<Stream> _innerStreamMock;
private readonly Mock<IProgress<StreamProgressArgs>> _progressMock;
private readonly ProgressStream _sut;
public ProgressStreamTests()
{
// Setup the mocks
_innerStreamMock = new Mock<Stream>();
_innerStreamMock.Setup(s => s.Length).Returns(1024L);
_progressMock = new Mock<IProgress<StreamProgressArgs>>();
// Inject mocks into the System Under Test
_sut = new ProgressStream(_innerStreamMock.Object, _progressMock.Object);
}
[Fact]
public async Task ReadAsync_Should_CallInnerStreamAndReportProgress()
{
// Arrange
var buffer = new byte[10];
_innerStreamMock
.Setup(s => s.ReadAsync(buffer, 0, buffer.Length, CancellationToken.None))
.Returns(Task.FromResult(5));
// Act
await _sut.ReadAsync(buffer, 0, buffer.Length);
// Assert - Inner Stream Read was called
_innerStreamMock.Verify(s => s.ReadAsync(buffer, 0, buffer.Length, CancellationToken.None), Times.Once);
// Assert - Progress Report was called with the correct byte count
_progressMock.Verify(p => p.Report(It.IsAny<StreamProgressArgs>()), Times.Once);
}
[Fact]
public async Task WriteAsync_Should_CallInnerStreamAndReportProgress()
{
// Arrange
var buffer = new byte[10];
_innerStreamMock
.Setup(s => s.WriteAsync(buffer, 0, buffer.Length, CancellationToken.None))
.Returns(Task.FromResult(5));
// Act
await _sut.WriteAsync(buffer, 0, buffer.Length);
// Assert - Inner Stream Write was called
_innerStreamMock.Verify(s => s.WriteAsync(buffer, 0, buffer.Length, CancellationToken.None), Times.Once);
// Assert - Progress Report was called with the correct byte count
_progressMock.Verify(p => p.Report(It.IsAny<StreamProgressArgs>()), Times.Once);
}
public void Dispose()
{
_sut.Dispose();
}
}
@@ -0,0 +1,46 @@
using Speckle.Sdk.Pipelines.Progress;
namespace Speckle.Sdk.Tests.Unit.Pipelines.Progress;
public class RenderedStreamProgressTests
{
[Theory]
[InlineData(1, "B", 1.0)]
[InlineData(1024, "B", 1.0)]
[InlineData(1024 + 1, "KB", 1.0 / 1024)]
[InlineData(1024 * 1024, "KB", 1.0 / 1024)]
[InlineData(1024 * 1024 + 1, "MB", 1.0 / (1024 * 1024))]
[InlineData(1024 * 1024 * 1024, "MB", 1.0 / (1024 * 1024))]
[InlineData(1024 * 1024 * 1024 + 1, "GB", 1.0 / (1024 * 1024 * 1024))]
[InlineData(1024L * 1024L * 1024L * 1024L, "GB", 1.0 / (1024L * 1024L * 1024L))]
public void GetFileSizeRendering_WithPositiveValue_ReturnsCorrectSuffix(
long value,
string expectedSuffix,
double expectedScaleFactor
)
{
var result = RenderedStreamProgress.GetFileSizeRendering(value);
Assert.Equal(expectedSuffix, result.suffix);
Assert.Equal(expectedScaleFactor, result.scaleFactor);
}
[Theory]
[InlineData(0)]
[InlineData(-1)]
[InlineData(-1000)]
public void GetFileSizeRendering_WithNonPositiveValue_ReturnsBytesSuffix(long value)
{
var result = RenderedStreamProgress.GetFileSizeRendering(value);
Assert.Equal("B", result.suffix);
Assert.Equal(1d, result.scaleFactor);
}
[Theory]
[InlineData(long.MaxValue)]
public void GetFileSizeRendering_WithVeryLargeValue_ThrowsArgumentOutOfRangeException(long value)
{
Assert.Throws<ArgumentOutOfRangeException>(() => RenderedStreamProgress.GetFileSizeRendering(value));
}
}