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
81 changes: 81 additions & 0 deletions GitLfsCache.Tests/Storage/FailingWriteStream.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
// Copyright (c) 2023-2026 ktsu-dev contributors

namespace ktsu.GitLfsCache.Tests.Storage;

/// <summary>
/// A sink that forwards a fixed number of writes and then fails every write after them the way a
/// full disk does.
/// </summary>
/// <param name="inner">The sink to forward the allowed writes to.</param>
/// <param name="allowedWrites">How many writes succeed before the first failure.</param>
internal sealed class FailingWriteStream(Stream inner, int allowedWrites) : Stream
{
private int _writes;

public override bool CanRead => false;

public override bool CanSeek => false;

public override bool CanWrite => true;

public override long Length => inner.Length;

public override long Position
{
get => inner.Position;
set
{
_ = value;
throw new NotSupportedException();
}
}

public override void Flush() => inner.Flush();

public override Task FlushAsync(CancellationToken cancellationToken) => inner.FlushAsync(cancellationToken);

public override int Read(byte[] buffer, int offset, int count) => throw new NotSupportedException();

public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException();

public override void SetLength(long value) => throw new NotSupportedException();

public override void Write(byte[] buffer, int offset, int count)
{
ThrowIfExhausted();
inner.Write(buffer, offset, count);
}

public override async ValueTask WriteAsync(ReadOnlyMemory<byte> buffer, CancellationToken cancellationToken = default)
{
ThrowIfExhausted();
await inner.WriteAsync(buffer, cancellationToken).ConfigureAwait(false);
}

public override Task WriteAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken) =>
WriteAsync(buffer.AsMemory(offset, count), cancellationToken).AsTask();

public override async ValueTask DisposeAsync()
{
await inner.DisposeAsync().ConfigureAwait(false);
await base.DisposeAsync().ConfigureAwait(false);
}

protected override void Dispose(bool disposing)
{
if (disposing)
{
inner.Dispose();
}

base.Dispose(disposing);
}

private void ThrowIfExhausted()
{
if (_writes++ >= allowedWrites)
{
throw new IOException("No space left on device");
}
}
}
38 changes: 38 additions & 0 deletions GitLfsCache.Tests/Storage/ObjectStoreTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@
using System.Text;
using ktsu.GitLfsCache.Configuration;
using ktsu.GitLfsCache.Storage;
using ktsu.Semantics.Paths;
using ktsu.Semantics.Strings;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
Expand Down Expand Up @@ -119,6 +121,42 @@
Assert.IsFalse(fileSystem.File.Exists(stagingPath), "Staging must not survive a mismatch.");
}

[TestMethod]
[DataRow(60_000, 0, DisplayName = "Single chunk, and it fails")]
[DataRow(100_000, 1, DisplayName = "Two chunks, and the last one fails")]
public async Task PublishAsync_StagingWriteFailedOnTheLastChunk_DoesNotPublish(int size, int allowedWrites)
{
(ObjectStore store, MockFileSystem fileSystem, _) = Create();
byte[] content = new byte[size];

for (int index = 0; index < size; index++)
{
content[index] = (byte)(index % 251);
}

string oid = Convert.ToHexStringLower(SHA256.HashData(content));

// The failing write is the last one the tee makes, so nothing after it could make the digest
// disagree with the oid. Only knowing that a write failed can stop the short file publishing.
string stagingPath = fileSystem.Path.Combine(Root, "staging.tmp");
Stream file = fileSystem.FileStream.New(stagingPath, FileMode.CreateNew, FileAccess.Write, FileShare.Read);
await using StagingHandle handle = new(
fileSystem,
stagingPath.As<AbsoluteFilePath>(),
new FailingWriteStream(file, allowedWrites));

using MemoryStream source = new(content);
using MemoryStream client = new();
await StreamTee.CopyAsync(source, client, handle.Stream, null, CancellationToken.None);

bool published = await store.PublishAsync(handle, "github", oid, CancellationToken.None);

CollectionAssert.AreEqual(content, client.ToArray(), "The client still gets every byte.");

Check warning on line 154 in GitLfsCache.Tests/Storage/ObjectStoreTests.cs

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Use 'Assert.AreSequenceEqual' instead of 'CollectionAssert.AreEqual'

See more on https://sonarcloud.io/project/issues?id=ktsu-dev_GitLfsCache&issues=AaDnYCPHUZvnceHz0yaS&open=AaDnYCPHUZvnceHz0yaS&pullRequest=59
Assert.IsFalse(published, "A staging file that missed a write is incomplete and must not publish.");
Assert.IsFalse(store.Exists("github", oid));
Assert.IsFalse(fileSystem.File.Exists(stagingPath), "Staging must not survive a failed write.");
}

[TestMethod]
public async Task PublishAsync_ObjectAlreadyPresent_SucceedsAndRemovesStaging()
{
Expand Down Expand Up @@ -249,7 +287,7 @@
}

[TestMethod]
public void Touch_MissingObject_DoesNotThrow()

Check warning on line 290 in GitLfsCache.Tests/Storage/ObjectStoreTests.cs

View workflow job for this annotation

GitHub Actions / Analyze & Release

Add at least one assertion to this test case.

Check warning on line 290 in GitLfsCache.Tests/Storage/ObjectStoreTests.cs

View workflow job for this annotation

GitHub Actions / Analyze & Release

Add at least one assertion to this test case.

Check warning on line 290 in GitLfsCache.Tests/Storage/ObjectStoreTests.cs

View workflow job for this annotation

GitHub Actions / Analyze & Release

Add at least one assertion to this test case.
{
(ObjectStore store, _, _) = Create();

Expand Down
33 changes: 31 additions & 2 deletions GitLfsCache/Storage/HashingStream.cs
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,15 @@ internal sealed class HashingStream(Stream inner, bool ownsInner = true) : Strea
/// <inheritdoc />
public override long Length => _written;

/// <summary>
/// Gets a value indicating whether a write to the inner sink has failed.
/// </summary>
/// <remarks>
/// Once a write fails the sink holds an unknown prefix of the content, so the digest no longer
/// describes what is on disk even when it matches. Whoever publishes the sink has to check this.
/// </remarks>
public bool Faulted { get; private set; }

/// <inheritdoc />
public override long Position
{
Expand Down Expand Up @@ -79,8 +88,19 @@ public override void Write(byte[] buffer, int offset, int count) =>
/// <inheritdoc />
public override void Write(ReadOnlySpan<byte> buffer)
{
try
{
inner.Write(buffer);
}
catch
{
Faulted = true;
throw;
}

// Digested only once the bytes have reached the sink, so a failed write cannot leave the
// digest describing content the sink never received.
_hash.AppendData(buffer);
inner.Write(buffer);
_written += buffer.Length;
}

Expand All @@ -89,8 +109,17 @@ public override async ValueTask WriteAsync(
ReadOnlyMemory<byte> buffer,
CancellationToken cancellationToken = default)
{
try
{
await inner.WriteAsync(buffer, cancellationToken).ConfigureAwait(false);
}
catch
{
Faulted = true;
throw;
}

_hash.AppendData(buffer.Span);
await inner.WriteAsync(buffer, cancellationToken).ConfigureAwait(false);
_written += buffer.Length;
}

Expand Down
10 changes: 10 additions & 0 deletions GitLfsCache/Storage/ObjectStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,16 @@ public async Task<bool> PublishAsync(
string digest = handle.GetDigestHex();
await handle.CloseAsync(cancellationToken).ConfigureAwait(false);

// A tee abandons a failed staging write and carries on serving the client, so a faulted handle
// can reach here. Its file is short by at least the failed write, and when that write was the
// last one the digest can still match, so the digest alone is not enough to publish on.
if (handle.Faulted)
{
StoreLog.DiscardedIncompleteObject(logger, upstream, oid);
await handle.DisposeAsync().ConfigureAwait(false);
return false;
}

if (!string.Equals(digest, oid, StringComparison.OrdinalIgnoreCase))
{
StoreLog.DiscardedMismatchedObject(logger, upstream, digest, oid);
Expand Down
6 changes: 6 additions & 0 deletions GitLfsCache/Storage/StagingHandle.cs
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,12 @@ internal StagingHandle(IFileSystem fileSystem, AbsoluteFilePath path, Stream sin
/// <returns>The digest, in the form a Git LFS object id takes.</returns>
public string GetDigestHex() => _stream.GetDigestHex();

/// <summary>
/// Gets a value indicating whether a write to the staging file failed, which leaves it incomplete
/// whatever its digest says.
/// </summary>
public bool Faulted => _stream.Faulted;

/// <summary>Marks the staging file as published so disposal leaves it alone.</summary>
internal void MarkPublished() => _published = true;

Expand Down
6 changes: 6 additions & 0 deletions GitLfsCache/Storage/StoreLog.cs
Original file line number Diff line number Diff line change
Expand Up @@ -79,4 +79,10 @@ public static partial void EvictedObjects(
Level = LogLevel.Error,
Message = "The store maintenance sweep failed and will be retried on the next interval.")]
public static partial void MaintenanceSweepFailed(ILogger logger, Exception exception);

[LoggerMessage(
EventId = 1009,
Level = LogLevel.Warning,
Message = "Discarding fetched object {Oid} for {Upstream}: a write to its staging file failed, so the file is incomplete.")]
public static partial void DiscardedIncompleteObject(ILogger logger, string upstream, string oid);
}
Loading