From 634cc7a3b5c28c8c4606ae15563c8913bf71086c Mon Sep 17 00:00:00 2001 From: Claude Date: Mon, 28 Sep 2026 09:24:37 +0000 Subject: [PATCH 1/2] Refuse to publish a staging file whose write failed HashingStream appended each chunk to the digest before writing it, so a write that failed left bytes in the digest that never reached disk. StreamTee abandons a failed staging sink and keeps serving the client, and the download path published regardless, so when the failed write was the last one the short file still matched its oid and was published. Every later hit served the truncated object. HashingStream now digests only after the inner write succeeds and records that a write failed. StagingHandle exposes that as Faulted, and ObjectStore.PublishAsync discards a faulted handle before comparing digests, which covers the download and upload paths alike. Fixes #45 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_016wsoxnwaqMuAzvkm2xnzYh --- .../Storage/FailingWriteStream.cs | 77 +++++++++++++++++++ GitLfsCache.Tests/Storage/ObjectStoreTests.cs | 38 +++++++++ GitLfsCache/Storage/HashingStream.cs | 33 +++++++- GitLfsCache/Storage/ObjectStore.cs | 10 +++ GitLfsCache/Storage/StagingHandle.cs | 6 ++ GitLfsCache/Storage/StoreLog.cs | 6 ++ 6 files changed, 168 insertions(+), 2 deletions(-) create mode 100644 GitLfsCache.Tests/Storage/FailingWriteStream.cs diff --git a/GitLfsCache.Tests/Storage/FailingWriteStream.cs b/GitLfsCache.Tests/Storage/FailingWriteStream.cs new file mode 100644 index 0000000..1f72291 --- /dev/null +++ b/GitLfsCache.Tests/Storage/FailingWriteStream.cs @@ -0,0 +1,77 @@ +// 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 => 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); } From da74169b49f020ee1b58319ca453df28fb3cbddb Mon Sep 17 00:00:00 2001 From: Claude Date: Mon, 28 Sep 2026 09:30:28 +0000 Subject: [PATCH 2/2] Discard the ignored Position value in FailingWriteStream Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_016wsoxnwaqMuAzvkm2xnzYh --- GitLfsCache.Tests/Storage/FailingWriteStream.cs | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/GitLfsCache.Tests/Storage/FailingWriteStream.cs b/GitLfsCache.Tests/Storage/FailingWriteStream.cs index 1f72291..69cf360 100644 --- a/GitLfsCache.Tests/Storage/FailingWriteStream.cs +++ b/GitLfsCache.Tests/Storage/FailingWriteStream.cs @@ -23,7 +23,11 @@ internal sealed class FailingWriteStream(Stream inner, int allowedWrites) : Stre public override long Position { get => inner.Position; - set => throw new NotSupportedException(); + set + { + _ = value; + throw new NotSupportedException(); + } } public override void Flush() => inner.Flush();