diff --git a/GitLfsCache.Tests/Storage/FailingWriteStream.cs b/GitLfsCache.Tests/Storage/FailingWriteStream.cs new file mode 100644 index 0000000..69cf360 --- /dev/null +++ b/GitLfsCache.Tests/Storage/FailingWriteStream.cs @@ -0,0 +1,81 @@ +// Copyright (c) 2023-2026 ktsu-dev contributors + +namespace ktsu.GitLfsCache.Tests.Storage; + +/// +/// A sink that forwards a fixed number of writes and then fails every write after them the way a +/// full disk does. +/// +/// The sink to forward the allowed writes to. +/// How many writes succeed before the first failure. +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 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"); + } + } +} diff --git a/GitLfsCache.Tests/Storage/ObjectStoreTests.cs b/GitLfsCache.Tests/Storage/ObjectStoreTests.cs index 8d8e94f..653d220 100644 --- a/GitLfsCache.Tests/Storage/ObjectStoreTests.cs +++ b/GitLfsCache.Tests/Storage/ObjectStoreTests.cs @@ -6,6 +6,8 @@ namespace ktsu.GitLfsCache.Tests.Storage; 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; @@ -119,6 +121,42 @@ public async Task PublishAsync_HashMismatch_DiscardsStagingAndReturnsFalse() 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(), + 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."); + 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() { diff --git a/GitLfsCache/Storage/HashingStream.cs b/GitLfsCache/Storage/HashingStream.cs index e524f27..d114416 100644 --- a/GitLfsCache/Storage/HashingStream.cs +++ b/GitLfsCache/Storage/HashingStream.cs @@ -38,6 +38,15 @@ internal sealed class HashingStream(Stream inner, bool ownsInner = true) : Strea /// public override long Length => _written; + /// + /// Gets a value indicating whether a write to the inner sink has failed. + /// + /// + /// 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. + /// + public bool Faulted { get; private set; } + /// public override long Position { @@ -79,8 +88,19 @@ public override void Write(byte[] buffer, int offset, int count) => /// public override void Write(ReadOnlySpan 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; } @@ -89,8 +109,17 @@ public override async ValueTask WriteAsync( ReadOnlyMemory 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; } diff --git a/GitLfsCache/Storage/ObjectStore.cs b/GitLfsCache/Storage/ObjectStore.cs index a892d41..2b1cf07 100644 --- a/GitLfsCache/Storage/ObjectStore.cs +++ b/GitLfsCache/Storage/ObjectStore.cs @@ -127,6 +127,16 @@ public async Task 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); diff --git a/GitLfsCache/Storage/StagingHandle.cs b/GitLfsCache/Storage/StagingHandle.cs index 958642c..5610c98 100644 --- a/GitLfsCache/Storage/StagingHandle.cs +++ b/GitLfsCache/Storage/StagingHandle.cs @@ -45,6 +45,12 @@ internal StagingHandle(IFileSystem fileSystem, AbsoluteFilePath path, Stream sin /// The digest, in the form a Git LFS object id takes. public string GetDigestHex() => _stream.GetDigestHex(); + /// + /// Gets a value indicating whether a write to the staging file failed, which leaves it incomplete + /// whatever its digest says. + /// + public bool Faulted => _stream.Faulted; + /// Marks the staging file as published so disposal leaves it alone. internal void MarkPublished() => _published = true; diff --git a/GitLfsCache/Storage/StoreLog.cs b/GitLfsCache/Storage/StoreLog.cs index 9a3a0a5..ae71bef 100644 --- a/GitLfsCache/Storage/StoreLog.cs +++ b/GitLfsCache/Storage/StoreLog.cs @@ -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); }