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);
}