From 460d4c3f17a6271416e7eef9b9f19ae23295640f Mon Sep 17 00:00:00 2001 From: Mirko Hirsch Date: Mon, 15 Dec 2025 15:04:38 +0100 Subject: [PATCH 1/7] add stream upload --- CHANGELOG.md | 8 ++++ .../BlockingContainerImageRegistryClient.kt | 5 ++- .../de/cmdjulian/kirc/client/UploadMode.kt | 8 ++++ .../ReactiveContainerImageRegistryClient.kt | 7 ++-- .../SuspendingContainerImageRegistryClient.kt | 16 ++++++-- .../kirc/impl/ContainerRegistryApi.kt | 9 ++++- .../kirc/impl/ContainerRegistryApiImpl.kt | 13 ++++--- ...pendingContainerImageRegistryClientImpl.kt | 34 ++++++++++++++-- .../kirc/impl/delegate/ImageUploader.kt | 17 ++++++-- .../de/cmdjulian/kirc/BlockingRegistryTest.kt | 39 +++++++++++++++++-- .../cmdjulian/kirc/DockerRegistryCliHelper.kt | 4 +- 11 files changed, 133 insertions(+), 27 deletions(-) create mode 100644 kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt diff --git a/CHANGELOG.md b/CHANGELOG.md index e47656e6..0df7160a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,14 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). ## [Unreleased] +### Added + +- Streamed upload + +### Changed + +- Upload modes available: Stream, Chunked, Compatibility (old slow version) + ## [v1.3.4] - 2025-12-10 ### Added diff --git a/kirc-blocking/src/main/kotlin/de/cmdjulian/kirc/client/BlockingContainerImageRegistryClient.kt b/kirc-blocking/src/main/kotlin/de/cmdjulian/kirc/client/BlockingContainerImageRegistryClient.kt index 76d91fe1..a5d36b0f 100644 --- a/kirc-blocking/src/main/kotlin/de/cmdjulian/kirc/client/BlockingContainerImageRegistryClient.kt +++ b/kirc-blocking/src/main/kotlin/de/cmdjulian/kirc/client/BlockingContainerImageRegistryClient.kt @@ -88,7 +88,7 @@ interface BlockingContainerImageRegistryClient { * * @return the digest of uploaded image */ - fun upload(repository: Repository, reference: Reference, tar: InputStream): Digest + fun upload(repository: Repository, reference: Reference, tar: InputStream, mode: UploadMode = UploadMode.Stream): Digest /** * Downloads a docker image for certain [reference]. @@ -144,12 +144,13 @@ fun SuspendingContainerImageRegistryClient.toBlockingClient() = object : Blockin return client.toBlockingClient() } - override fun upload(repository: Repository, reference: Reference, tar: InputStream): Digest = + override fun upload(repository: Repository, reference: Reference, tar: InputStream, mode: UploadMode): Digest = runBlocking(Dispatchers.Default) { this@toBlockingClient.upload( repository, reference, tar.asSource().buffered(), + mode ) } diff --git a/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt b/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt new file mode 100644 index 00000000..4cea142a --- /dev/null +++ b/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt @@ -0,0 +1,8 @@ +package de.cmdjulian.kirc.client + +sealed class UploadMode { + data object Stream: UploadMode() + @Deprecated("Use Chunked and specify the chunk size") + data object Compatibility: UploadMode() + data class Chunked(val chunkSize: Long): UploadMode() +} \ No newline at end of file diff --git a/kirc-reactive/src/main/kotlin/de/cmdjulian/kirc/client/ReactiveContainerImageRegistryClient.kt b/kirc-reactive/src/main/kotlin/de/cmdjulian/kirc/client/ReactiveContainerImageRegistryClient.kt index 22a1ffa1..82d8f7fb 100644 --- a/kirc-reactive/src/main/kotlin/de/cmdjulian/kirc/client/ReactiveContainerImageRegistryClient.kt +++ b/kirc-reactive/src/main/kotlin/de/cmdjulian/kirc/client/ReactiveContainerImageRegistryClient.kt @@ -4,6 +4,7 @@ import de.cmdjulian.kirc.image.Digest import de.cmdjulian.kirc.image.Reference import de.cmdjulian.kirc.image.Repository import de.cmdjulian.kirc.image.Tag + import de.cmdjulian.kirc.spec.image.ImageConfig import de.cmdjulian.kirc.spec.manifest.Manifest import de.cmdjulian.kirc.spec.manifest.ManifestList @@ -89,7 +90,7 @@ interface ReactiveContainerImageRegistryClient { * * @return the digest of uploaded image */ - fun upload(repository: Repository, reference: Reference, tar: Flux): Mono + fun upload(repository: Repository, reference: Reference, tar: Flux, mode: UploadMode = UploadMode.Stream): Mono /** * Downloads a docker image for certain [reference]. @@ -141,9 +142,9 @@ fun SuspendingContainerImageRegistryClient.toReactiveClient() = object : Reactiv } } - override fun upload(repository: Repository, reference: Reference, tar: Flux): Mono = mono { + override fun upload(repository: Repository, reference: Reference, tar: Flux, mode: UploadMode): Mono = mono { val buffer = Buffer().also { buffer -> tar.collect(buffer::writeByte) } - this@toReactiveClient.upload(repository, reference, buffer) + this@toReactiveClient.upload(repository, reference, buffer, mode) } override fun download(repository: Repository, reference: Reference): Flux = flux { diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/client/SuspendingContainerImageRegistryClient.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/client/SuspendingContainerImageRegistryClient.kt index 49cd620b..e87c6be0 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/client/SuspendingContainerImageRegistryClient.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/client/SuspendingContainerImageRegistryClient.kt @@ -4,8 +4,10 @@ import de.cmdjulian.kirc.image.Digest import de.cmdjulian.kirc.image.Reference import de.cmdjulian.kirc.image.Repository import de.cmdjulian.kirc.image.Tag + import de.cmdjulian.kirc.impl.response.ResultSource import de.cmdjulian.kirc.impl.response.UploadSession +import de.cmdjulian.kirc.spec.UploadBlobPath import de.cmdjulian.kirc.spec.image.ImageConfig import de.cmdjulian.kirc.spec.manifest.Manifest import de.cmdjulian.kirc.spec.manifest.ManifestList @@ -113,9 +115,17 @@ interface SuspendingContainerImageRegistryClient { suspend fun uploadBlobChunks(session: UploadSession, path: Path, chunkSize: Long = 10 * 1048576L): UploadSession /** - * Uploads an entire blob by stream + * Uploads an entire blob chunk-wise for reduced memory load + */ + @Deprecated("Use chunked upload instead") + suspend fun uploadBlobCompatibility(session: UploadSession, path: Path): UploadSession + + /** + * Uploads an entire [blob] by stream + * + * Closes [session] by uploading the whole blob in one monolithic upload and returns its digest upon success. */ - suspend fun uploadBlobStream(session: UploadSession, stream: Source): UploadSession + suspend fun uploadBlobStream(session: UploadSession, blob: UploadBlobPath): Digest /** * Upload a manifest @@ -140,7 +150,7 @@ interface SuspendingContainerImageRegistryClient { * * @return the digest of uploaded image */ - suspend fun upload(repository: Repository, reference: Reference, tar: Source): Digest + suspend fun upload(repository: Repository, reference: Reference, tar: Source, mode: UploadMode = UploadMode.Stream): Digest /** * Downloads a docker image for certain [reference]. diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApi.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApi.kt index 6f25441c..61f0c962 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApi.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApi.kt @@ -13,6 +13,7 @@ import de.cmdjulian.kirc.spec.manifest.Manifest import de.cmdjulian.kirc.spec.manifest.ManifestSingle import kotlinx.io.Buffer import kotlinx.io.Source +import java.nio.file.Path /** * Defines the calls to the container registry API @@ -77,8 +78,12 @@ internal interface ContainerRegistryApi { endRange: Long, ): Result - /** Uploads the whole blob data [Source] as stream */ - suspend fun uploadBlobStream(session: UploadSession, source: Source): Result + /** + * Uploads the whole blob data from [path] with [size] and its [digest] as stream + * + * Upon success, should return the same [digest] and closes [session] + */ + suspend fun uploadBlobStream(session: UploadSession, path: Path, size: Long, digest: Digest): Result /** Retrieve the status of provided [session], returning the range of already uploaded data (start, end) */ suspend fun uploadStatus(session: UploadSession): Result, FuelError> diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt index 7ccd660e..9f8f80a5 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt @@ -3,6 +3,7 @@ package de.cmdjulian.kirc.impl import com.github.kittinunf.fuel.core.FuelError import com.github.kittinunf.fuel.core.FuelManager import com.github.kittinunf.fuel.core.Headers +import com.github.kittinunf.fuel.core.Method import com.github.kittinunf.fuel.core.Parameters import com.github.kittinunf.fuel.core.awaitResponseResult import com.github.kittinunf.fuel.core.deserializers.ByteArrayDeserializer @@ -38,6 +39,9 @@ import kotlinx.io.Buffer import kotlinx.io.Source import kotlinx.io.asInputStream import kotlinx.io.readByteArray +import java.nio.file.Path +import kotlin.io.path.fileSize +import kotlin.io.path.inputStream private const val APPLICATION_JSON = "application/json" private const val APPLICATION_OCTET_STREAM = "application/octet-stream" @@ -263,14 +267,13 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .mapToUploadSession() - // Currently not working as intended, because internal fuel buffer has an overflow - override suspend fun uploadBlobStream(session: UploadSession, source: Source): Result = - fuelManager.patch(session.location) + override suspend fun uploadBlobStream(session: UploadSession, path: Path, size: Long, digest: Digest): Result = + fuelManager.put(session.location, listOf("digest" to digest)) .appendHeader(Headers.CONTENT_TYPE, APPLICATION_OCTET_STREAM) - .body(source::asInputStream) + .body(path::inputStream, { size }, repeatable = true) .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } - .mapToUploadSession() + .mapToDigest() override suspend fun uploadStatus(session: UploadSession): Result, FuelError> = fuelManager.get(session.location) diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/SuspendingContainerImageRegistryClientImpl.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/SuspendingContainerImageRegistryClientImpl.kt index d950aae9..823d808d 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/SuspendingContainerImageRegistryClientImpl.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/SuspendingContainerImageRegistryClientImpl.kt @@ -5,6 +5,7 @@ import com.github.kittinunf.result.map import com.github.kittinunf.result.onError import de.cmdjulian.kirc.client.SuspendingContainerImageClient import de.cmdjulian.kirc.client.SuspendingContainerImageRegistryClient +import de.cmdjulian.kirc.client.UploadMode import de.cmdjulian.kirc.image.ContainerImageName import de.cmdjulian.kirc.image.Digest import de.cmdjulian.kirc.image.Reference @@ -12,10 +13,12 @@ import de.cmdjulian.kirc.image.Repository import de.cmdjulian.kirc.image.Tag import de.cmdjulian.kirc.impl.delegate.ImageDownloader import de.cmdjulian.kirc.impl.delegate.ImageUploader + import de.cmdjulian.kirc.impl.response.Catalog import de.cmdjulian.kirc.impl.response.ResultSource import de.cmdjulian.kirc.impl.response.TagList import de.cmdjulian.kirc.impl.response.UploadSession +import de.cmdjulian.kirc.spec.UploadBlobPath import de.cmdjulian.kirc.spec.image.DockerImageConfigV1 import de.cmdjulian.kirc.spec.image.ImageConfig import de.cmdjulian.kirc.spec.image.OciImageConfigV1 @@ -34,6 +37,7 @@ import kotlinx.io.Source import kotlinx.io.asInputStream import kotlinx.io.buffered import kotlinx.io.files.SystemFileSystem +import kotlinx.io.readAtMostTo import java.nio.file.Path internal class SuspendingContainerImageRegistryClientImpl(private val api: ContainerRegistryApi, tmpPath: Path) : @@ -124,8 +128,8 @@ internal class SuspendingContainerImageRegistryClientImpl(private val api: Conta override suspend fun initiateBlobUpload(repository: Repository): UploadSession = api.initiateUpload(repository).getOrElse { throw it.toRegistryClientError(repository, null) } - override suspend fun uploadBlobStream(session: UploadSession, stream: Source): UploadSession = - api.uploadBlobStream(session, stream).getOrElse { throw it.toRegistryClientError() } + override suspend fun uploadBlobStream(session: UploadSession, blob: UploadBlobPath): Digest = + api.uploadBlobStream(session, blob.path, blob.size, blob.digest).getOrElse { throw it.toRegistryClientError() } override suspend fun uploadBlobChunks(session: UploadSession, path: Path, chunkSize: Long): UploadSession = withContext(Dispatchers.IO) { SystemFileSystem.source(path.toKotlinPath()).buffered() }.use { stream -> @@ -151,6 +155,28 @@ internal class SuspendingContainerImageRegistryClientImpl(private val api: Conta currentSession } + @Deprecated("Use chunked upload instead") + override suspend fun uploadBlobCompatibility(session: UploadSession, path: Path): UploadSession = + withContext(Dispatchers.IO) { + SystemFileSystem.source(path.toKotlinPath()).buffered() + }.use { stream -> + var returnedSession = session + var startRange = 0L + var endRange: Long + + while (!stream.exhausted()) { + val buffer = Buffer() + // will not read more than 8KB into buffer but this version worked always + stream.readAtMostTo(buffer, 10 * 1048576L) + endRange = startRange + buffer.size - 1 + returnedSession = api.uploadBlobChunked(returnedSession, buffer, startRange, endRange) + .getOrElse { throw it.toRegistryClientError() } + startRange = endRange + } + + returnedSession + } + override suspend fun finishBlobUpload(session: UploadSession, digest: Digest): Digest = api.finishBlobUpload(session, digest) .getOrElse { throw it.toRegistryClientError(null, digest) } @@ -166,8 +192,8 @@ internal class SuspendingContainerImageRegistryClientImpl(private val api: Conta api.uploadManifest(repository, reference, manifest) .getOrElse { throw it.toRegistryClientError(repository, reference) } - override suspend fun upload(repository: Repository, reference: Reference, tar: Source): Digest = - uploader.upload(repository, reference, tar) + override suspend fun upload(repository: Repository, reference: Reference, tar: Source, mode: UploadMode): Digest = + uploader.upload(repository, reference, tar, mode) override suspend fun download(repository: Repository, reference: Reference): Source = downloader.download(repository, reference) diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/delegate/ImageUploader.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/delegate/ImageUploader.kt index cbeb16f7..4151e304 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/delegate/ImageUploader.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/delegate/ImageUploader.kt @@ -2,6 +2,7 @@ package de.cmdjulian.kirc.impl.delegate import de.cmdjulian.kirc.KircUploadException import de.cmdjulian.kirc.client.SuspendingContainerImageRegistryClient +import de.cmdjulian.kirc.client.UploadMode import de.cmdjulian.kirc.image.Digest import de.cmdjulian.kirc.image.Reference import de.cmdjulian.kirc.image.Repository @@ -36,7 +37,7 @@ import kotlin.io.path.pathString internal class ImageUploader(private val client: SuspendingContainerImageRegistryClient, private val tmpPath: Path) { - suspend fun upload(repository: Repository, reference: Reference, tar: Source): Digest = coroutineScope { + suspend fun upload(repository: Repository, reference: Reference, tar: Source, upload: UploadMode = UploadMode.Stream): Digest = coroutineScope { // store data temporarily val tempDirectory = Path.of( tmpPath.pathString, @@ -62,9 +63,17 @@ internal class ImageUploader(private val client: SuspendingContainerImageRegistr if (!client.existsBlob(repository, blob.digest)) { val session = client.initiateBlobUpload(repository) - val endSession = client.uploadBlobChunks(session, blob.path) - - client.finishBlobUpload(endSession, blob.digest) + when(upload) { + is UploadMode.Chunked -> { + val endSession = client.uploadBlobChunks(session, blob.path, upload.chunkSize) + client.finishBlobUpload(endSession, blob.digest) + } + is UploadMode.Compatibility -> { + val endSession = client.uploadBlobCompatibility(session, blob.path) + client.finishBlobUpload(endSession, blob.digest) + } + is UploadMode.Stream -> client.uploadBlobStream(session, blob) + } } } } diff --git a/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/BlockingRegistryTest.kt b/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/BlockingRegistryTest.kt index da283e6f..b62aabf5 100644 --- a/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/BlockingRegistryTest.kt +++ b/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/BlockingRegistryTest.kt @@ -3,9 +3,11 @@ package de.cmdjulian.kirc import de.cmdjulian.kirc.client.BlockingContainerImageClientFactory import de.cmdjulian.kirc.client.BlockingContainerImageRegistryClient import de.cmdjulian.kirc.client.RegistryCredentials +import de.cmdjulian.kirc.client.UploadMode import de.cmdjulian.kirc.image.Digest import de.cmdjulian.kirc.image.Repository import de.cmdjulian.kirc.image.Tag + import de.cmdjulian.kirc.spec.manifest.ManifestList import io.kotest.assertions.throwables.shouldNotThrowAny import io.kotest.matchers.booleans.shouldBeTrue @@ -151,7 +153,38 @@ internal class BlockingRegistryTest { } @Test - fun `upload - to registry`() { + fun `upload stream - to registry`() { + val data = SystemFileSystem.source(Path(helloWorldImage.path)) + val repository = Repository("python") + val tag = Tag("test") + + client.exists(repository, tag) shouldBe false + + shouldNotThrowAny { + client.upload(repository, tag, data.buffered().asInputStream(), UploadMode.Stream) + } + + client.exists(repository, tag) shouldBe true + } + + @Test + fun `upload chunked - to registry`() { + val data = SystemFileSystem.source(Path(helloWorldImage.path)) + val repository = Repository("python") + val tag = Tag("test") + val upload = UploadMode.Chunked(10 * 1048576L) // 10MB + + client.exists(repository, tag) shouldBe false + + shouldNotThrowAny { + client.upload(repository, tag, data.buffered().asInputStream(), upload) + } + + client.exists(repository, tag) shouldBe true + } + + @Test + fun `upload compatibility - to registry`() { val data = SystemFileSystem.source(Path(helloWorldImage.path)) val repository = Repository("python") val tag = Tag("test") @@ -159,7 +192,7 @@ internal class BlockingRegistryTest { client.exists(repository, tag) shouldBe false shouldNotThrowAny { - client.upload(repository, tag, data.buffered().asInputStream()) + client.upload(repository, tag, data.buffered().asInputStream(), UploadMode.Compatibility) } client.exists(repository, tag) shouldBe true @@ -175,7 +208,7 @@ internal class BlockingRegistryTest { val result = client.download(repository, tag) shouldNotThrowAny { // check if upload of downloaded data possible - client.upload(Repository("test"), Tag("upload"), result) + client.upload(Repository("test"), Tag("upload"), result, UploadMode.Stream) } } diff --git a/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/DockerRegistryCliHelper.kt b/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/DockerRegistryCliHelper.kt index af1cf3bb..39d0ec12 100644 --- a/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/DockerRegistryCliHelper.kt +++ b/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/DockerRegistryCliHelper.kt @@ -3,9 +3,11 @@ package de.cmdjulian.kirc import de.cmdjulian.kirc.client.BlockingContainerImageClientFactory import de.cmdjulian.kirc.client.BlockingContainerImageRegistryClient import de.cmdjulian.kirc.client.RegistryCredentials +import de.cmdjulian.kirc.client.UploadMode import de.cmdjulian.kirc.image.Digest import de.cmdjulian.kirc.image.Reference import de.cmdjulian.kirc.image.Repository + import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.runBlocking import kotlinx.io.asInputStream @@ -26,7 +28,7 @@ class DockerRegistryCliHelper(addressName: String, credentials: RegistryCredenti fun pushImage(repository: Repository, reference: Reference, url: URL): Digest { images.add(UploadReference(repository, reference)) val source = runBlocking(Dispatchers.IO) { SystemFileSystem.source(Path(url.path)) } - return client.upload(repository, reference, source.buffered().asInputStream()) + return client.upload(repository, reference, source.buffered().asInputStream(), UploadMode.Stream) } fun deleteAll() { From 5733e0f6c20d79c453442243f69ba624b2d95c42 Mon Sep 17 00:00:00 2001 From: Mirko Hirsch Date: Mon, 15 Dec 2025 15:11:22 +0100 Subject: [PATCH 2/7] ktlint refactor --- .../BlockingContainerImageRegistryClient.kt | 9 +++++++-- .../de/cmdjulian/kirc/client/UploadMode.kt | 7 ++++--- .../ReactiveContainerImageRegistryClient.kt | 16 +++++++++++----- .../SuspendingContainerImageRegistryClient.kt | 7 ++++++- .../cmdjulian/kirc/impl/ContainerRegistryApi.kt | 7 ++++++- .../kirc/impl/ContainerRegistryApiImpl.kt | 10 ++++++---- ...SuspendingContainerImageRegistryClientImpl.kt | 2 -- .../kirc/impl/delegate/ImageUploader.kt | 14 +++++++++++--- .../de/cmdjulian/kirc/utils/UtilityFuns.kt | 2 +- 9 files changed, 52 insertions(+), 22 deletions(-) diff --git a/kirc-blocking/src/main/kotlin/de/cmdjulian/kirc/client/BlockingContainerImageRegistryClient.kt b/kirc-blocking/src/main/kotlin/de/cmdjulian/kirc/client/BlockingContainerImageRegistryClient.kt index a5d36b0f..9ee9a7bf 100644 --- a/kirc-blocking/src/main/kotlin/de/cmdjulian/kirc/client/BlockingContainerImageRegistryClient.kt +++ b/kirc-blocking/src/main/kotlin/de/cmdjulian/kirc/client/BlockingContainerImageRegistryClient.kt @@ -88,7 +88,12 @@ interface BlockingContainerImageRegistryClient { * * @return the digest of uploaded image */ - fun upload(repository: Repository, reference: Reference, tar: InputStream, mode: UploadMode = UploadMode.Stream): Digest + fun upload( + repository: Repository, + reference: Reference, + tar: InputStream, + mode: UploadMode = UploadMode.Stream, + ): Digest /** * Downloads a docker image for certain [reference]. @@ -150,7 +155,7 @@ fun SuspendingContainerImageRegistryClient.toBlockingClient() = object : Blockin repository, reference, tar.asSource().buffered(), - mode + mode, ) } diff --git a/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt b/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt index 4cea142a..5a6ec82b 100644 --- a/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt +++ b/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt @@ -1,8 +1,9 @@ package de.cmdjulian.kirc.client sealed class UploadMode { - data object Stream: UploadMode() + data object Stream : UploadMode() + @Deprecated("Use Chunked and specify the chunk size") - data object Compatibility: UploadMode() - data class Chunked(val chunkSize: Long): UploadMode() + data object Compatibility : UploadMode() + data class Chunked(val chunkSize: Long) : UploadMode() } \ No newline at end of file diff --git a/kirc-reactive/src/main/kotlin/de/cmdjulian/kirc/client/ReactiveContainerImageRegistryClient.kt b/kirc-reactive/src/main/kotlin/de/cmdjulian/kirc/client/ReactiveContainerImageRegistryClient.kt index 82d8f7fb..cc9b4686 100644 --- a/kirc-reactive/src/main/kotlin/de/cmdjulian/kirc/client/ReactiveContainerImageRegistryClient.kt +++ b/kirc-reactive/src/main/kotlin/de/cmdjulian/kirc/client/ReactiveContainerImageRegistryClient.kt @@ -90,7 +90,12 @@ interface ReactiveContainerImageRegistryClient { * * @return the digest of uploaded image */ - fun upload(repository: Repository, reference: Reference, tar: Flux, mode: UploadMode = UploadMode.Stream): Mono + fun upload( + repository: Repository, + reference: Reference, + tar: Flux, + mode: UploadMode = UploadMode.Stream, + ): Mono /** * Downloads a docker image for certain [reference]. @@ -142,10 +147,11 @@ fun SuspendingContainerImageRegistryClient.toReactiveClient() = object : Reactiv } } - override fun upload(repository: Repository, reference: Reference, tar: Flux, mode: UploadMode): Mono = mono { - val buffer = Buffer().also { buffer -> tar.collect(buffer::writeByte) } - this@toReactiveClient.upload(repository, reference, buffer, mode) - } + override fun upload(repository: Repository, reference: Reference, tar: Flux, mode: UploadMode): Mono = + mono { + val buffer = Buffer().also { buffer -> tar.collect(buffer::writeByte) } + this@toReactiveClient.upload(repository, reference, buffer, mode) + } override fun download(repository: Repository, reference: Reference): Flux = flux { this@toReactiveClient.download(repository, reference).use { result -> diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/client/SuspendingContainerImageRegistryClient.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/client/SuspendingContainerImageRegistryClient.kt index e87c6be0..80ba5322 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/client/SuspendingContainerImageRegistryClient.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/client/SuspendingContainerImageRegistryClient.kt @@ -150,7 +150,12 @@ interface SuspendingContainerImageRegistryClient { * * @return the digest of uploaded image */ - suspend fun upload(repository: Repository, reference: Reference, tar: Source, mode: UploadMode = UploadMode.Stream): Digest + suspend fun upload( + repository: Repository, + reference: Reference, + tar: Source, + mode: UploadMode = UploadMode.Stream, + ): Digest /** * Downloads a docker image for certain [reference]. diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApi.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApi.kt index 61f0c962..633bb5ec 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApi.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApi.kt @@ -83,7 +83,12 @@ internal interface ContainerRegistryApi { * * Upon success, should return the same [digest] and closes [session] */ - suspend fun uploadBlobStream(session: UploadSession, path: Path, size: Long, digest: Digest): Result + suspend fun uploadBlobStream( + session: UploadSession, + path: Path, + size: Long, + digest: Digest, + ): Result /** Retrieve the status of provided [session], returning the range of already uploaded data (start, end) */ suspend fun uploadStatus(session: UploadSession): Result, FuelError> diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt index 9f8f80a5..4f656b33 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt @@ -3,7 +3,6 @@ package de.cmdjulian.kirc.impl import com.github.kittinunf.fuel.core.FuelError import com.github.kittinunf.fuel.core.FuelManager import com.github.kittinunf.fuel.core.Headers -import com.github.kittinunf.fuel.core.Method import com.github.kittinunf.fuel.core.Parameters import com.github.kittinunf.fuel.core.awaitResponseResult import com.github.kittinunf.fuel.core.deserializers.ByteArrayDeserializer @@ -37,10 +36,8 @@ import de.cmdjulian.kirc.utils.mapToResultSource import de.cmdjulian.kirc.utils.mapToUploadSession import kotlinx.io.Buffer import kotlinx.io.Source -import kotlinx.io.asInputStream import kotlinx.io.readByteArray import java.nio.file.Path -import kotlin.io.path.fileSize import kotlin.io.path.inputStream private const val APPLICATION_JSON = "application/json" @@ -267,7 +264,12 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .mapToUploadSession() - override suspend fun uploadBlobStream(session: UploadSession, path: Path, size: Long, digest: Digest): Result = + override suspend fun uploadBlobStream( + session: UploadSession, + path: Path, + size: Long, + digest: Digest, + ): Result = fuelManager.put(session.location, listOf("digest" to digest)) .appendHeader(Headers.CONTENT_TYPE, APPLICATION_OCTET_STREAM) .body(path::inputStream, { size }, repeatable = true) diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/SuspendingContainerImageRegistryClientImpl.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/SuspendingContainerImageRegistryClientImpl.kt index 823d808d..897c0c84 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/SuspendingContainerImageRegistryClientImpl.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/SuspendingContainerImageRegistryClientImpl.kt @@ -13,7 +13,6 @@ import de.cmdjulian.kirc.image.Repository import de.cmdjulian.kirc.image.Tag import de.cmdjulian.kirc.impl.delegate.ImageDownloader import de.cmdjulian.kirc.impl.delegate.ImageUploader - import de.cmdjulian.kirc.impl.response.Catalog import de.cmdjulian.kirc.impl.response.ResultSource import de.cmdjulian.kirc.impl.response.TagList @@ -37,7 +36,6 @@ import kotlinx.io.Source import kotlinx.io.asInputStream import kotlinx.io.buffered import kotlinx.io.files.SystemFileSystem -import kotlinx.io.readAtMostTo import java.nio.file.Path internal class SuspendingContainerImageRegistryClientImpl(private val api: ContainerRegistryApi, tmpPath: Path) : diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/delegate/ImageUploader.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/delegate/ImageUploader.kt index 4151e304..6f838019 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/delegate/ImageUploader.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/delegate/ImageUploader.kt @@ -37,7 +37,12 @@ import kotlin.io.path.pathString internal class ImageUploader(private val client: SuspendingContainerImageRegistryClient, private val tmpPath: Path) { - suspend fun upload(repository: Repository, reference: Reference, tar: Source, upload: UploadMode = UploadMode.Stream): Digest = coroutineScope { + suspend fun upload( + repository: Repository, + reference: Reference, + tar: Source, + upload: UploadMode = UploadMode.Stream, + ): Digest = coroutineScope { // store data temporarily val tempDirectory = Path.of( tmpPath.pathString, @@ -63,15 +68,18 @@ internal class ImageUploader(private val client: SuspendingContainerImageRegistr if (!client.existsBlob(repository, blob.digest)) { val session = client.initiateBlobUpload(repository) - when(upload) { + when (upload) { is UploadMode.Chunked -> { - val endSession = client.uploadBlobChunks(session, blob.path, upload.chunkSize) + val endSession = + client.uploadBlobChunks(session, blob.path, upload.chunkSize) client.finishBlobUpload(endSession, blob.digest) } + is UploadMode.Compatibility -> { val endSession = client.uploadBlobCompatibility(session, blob.path) client.finishBlobUpload(endSession, blob.digest) } + is UploadMode.Stream -> client.uploadBlobStream(session, blob) } } diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/utils/UtilityFuns.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/utils/UtilityFuns.kt index 58d005c9..35b94a76 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/utils/UtilityFuns.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/utils/UtilityFuns.kt @@ -1,7 +1,7 @@ package de.cmdjulian.kirc.utils import kotlin.io.path.pathString -import java.nio.file.Path as JavaPath import kotlinx.io.files.Path as KotlinPath +import java.nio.file.Path as JavaPath internal fun JavaPath.toKotlinPath(): KotlinPath = KotlinPath(pathString) From 03ab242181561ae2254c66b313d83a86faa79285 Mon Sep 17 00:00:00 2001 From: Mirko Hirsch Date: Mon, 15 Dec 2025 16:11:29 +0100 Subject: [PATCH 3/7] ktlint refactor --- .../main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt b/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt index 5a6ec82b..05e13157 100644 --- a/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt +++ b/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt @@ -1,9 +1,18 @@ package de.cmdjulian.kirc.client +/** + * Determines the way image blobs are uploaded. + * + * Currently supported modes: + * - [Stream]: uploads blob as stream in one request + * - [Chunked]: splits blob into [Chunked.chunkSize] bytes and uploads them one after another + * - [Compatibility]: uploads blob split into 8KB chunks. This is very slow but in some cases works best. Use Chunked instead. + */ sealed class UploadMode { data object Stream : UploadMode() @Deprecated("Use Chunked and specify the chunk size") data object Compatibility : UploadMode() + data class Chunked(val chunkSize: Long) : UploadMode() } \ No newline at end of file From b2bf73938f54b728f4645f1e4746ad34cc80f14e Mon Sep 17 00:00:00 2001 From: Mirko Hirsch Date: Mon, 15 Dec 2025 16:34:22 +0100 Subject: [PATCH 4/7] fix response retry --- CHANGELOG.md | 5 ++ .../kirc/impl/ContainerRegistryApiImpl.kt | 49 ++++++++++++++++--- .../impl/ResponseRetryWithAuthentication.kt | 17 +++++++ 3 files changed, 63 insertions(+), 8 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0df7160a..4ca3c9c1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -14,6 +14,11 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). - Upload modes available: Stream, Chunked, Compatibility (old slow version) +### Fixed + +- Requests only retry upon bearer authentication +- Streamed upload isn't stored in memory when request retries + ## [v1.3.4] - 2025-12-10 ### Added diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt index 4f656b33..86df81ba 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt @@ -4,9 +4,11 @@ import com.github.kittinunf.fuel.core.FuelError import com.github.kittinunf.fuel.core.FuelManager import com.github.kittinunf.fuel.core.Headers import com.github.kittinunf.fuel.core.Parameters +import com.github.kittinunf.fuel.core.Request import com.github.kittinunf.fuel.core.awaitResponseResult import com.github.kittinunf.fuel.core.deserializers.ByteArrayDeserializer import com.github.kittinunf.fuel.core.deserializers.EmptyDeserializer +import com.github.kittinunf.fuel.core.extensions.authentication import com.github.kittinunf.result.Result import com.github.kittinunf.result.flatMap import com.github.kittinunf.result.map @@ -43,14 +45,22 @@ import kotlin.io.path.inputStream private const val APPLICATION_JSON = "application/json" private const val APPLICATION_OCTET_STREAM = "application/octet-stream" -internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, credentials: RegistryCredentials?) : +internal class ContainerRegistryApiImpl( + private val fuelManager: FuelManager, + private val credentials: RegistryCredentials?, +) : ContainerRegistryApi { private val handler = ResponseRetryWithAuthentication(credentials, fuelManager) + private fun Request.addBasicAuth() = apply { + credentials?.let { credentials -> authentication().basic(credentials.username, credentials.password) } + } + // Status override suspend fun ping(): Result<*, FuelError> = fuelManager.get("/v2/") + .addBasicAuth() .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .third @@ -63,6 +73,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr val deserializable = jacksonDeserializer() return fuelManager.get("/v2/_catalog", parameter) + .addBasicAuth() .appendHeader(Headers.ACCEPT, APPLICATION_JSON) .awaitResponseResult(deserializable) .let { responseResult -> handler.retryOnUnauthorized(responseResult, deserializable) } @@ -77,6 +88,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr val deserializable = jacksonDeserializer() return fuelManager.get("/v2/$repository/tags/list", parameter) + .addBasicAuth() .appendHeader(Headers.ACCEPT, APPLICATION_JSON) .awaitResponseResult(deserializable) .let { responseResult -> handler.retryOnUnauthorized(responseResult, deserializable) } @@ -86,6 +98,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr override suspend fun digest(repository: Repository, reference: Reference) = when (reference) { is Digest -> Result.success(reference) is Tag -> fuelManager.head("/v2/$repository/manifests/$reference") + .addBasicAuth() .appendHeader( Headers.ACCEPT, APPLICATION_JSON, @@ -106,6 +119,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr reference: Reference, accept: String, ): Result<*, FuelError> = fuelManager.head("/v2/$repository/manifests/$reference") + .addBasicAuth() .appendHeader(Headers.ACCEPT, accept) .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } @@ -114,6 +128,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr override suspend fun manifests(repository: Repository, reference: Reference): Result { val deserializable = jacksonDeserializer() return fuelManager.get("/v2/$repository/manifests/$reference") + .addBasicAuth() .appendHeader( Headers.ACCEPT, APPLICATION_JSON, @@ -130,6 +145,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr override suspend fun manifestStream(repository: Repository, reference: Reference): Result { val deserializable = SourceDeserializer() return fuelManager.get("/v2/$repository/manifests/$reference") + .addBasicAuth() .appendHeader( Headers.ACCEPT, APPLICATION_JSON, @@ -146,6 +162,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr override suspend fun manifest(repository: Repository, reference: Reference): Result { val deserializable = jacksonDeserializer() return fuelManager.get("/v2/$repository/manifests/$reference") + .addBasicAuth() .appendHeader(Headers.ACCEPT, APPLICATION_JSON, OciManifestV1.MediaType, DockerManifestV2.MediaType) .awaitResponseResult(deserializable) .let { responseResult -> handler.retryOnUnauthorized(responseResult, deserializable) } @@ -168,6 +185,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr is OciManifestListV1 -> OciManifestListV1.MediaType } return fuelManager.put("/v2/$repository/manifests/$urlReference") + .addBasicAuth() .appendHeader(Headers.CONTENT_TYPE, contentType) .body(JsonMapper.writeValueAsString(manifest)) .awaitResponseResult(EmptyDeserializer) @@ -178,6 +196,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr override suspend fun deleteManifest(repository: Repository, reference: Reference): Result { return digest(repository, reference).flatMap { digest -> fuelManager.delete("/v2/$repository/manifests/$digest") + .addBasicAuth() .appendHeader( Headers.ACCEPT, APPLICATION_JSON, @@ -197,6 +216,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr override suspend fun existsBlob(repository: Repository, digest: Digest): Result<*, FuelError> = fuelManager.head("/v2/$repository/blobs/$digest") + .addBasicAuth() .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .third @@ -204,6 +224,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr override suspend fun blob(repository: Repository, digest: Digest): Result { val deserializable = ByteArrayDeserializer() return fuelManager.get("/v2/$repository/blobs/$digest") + .addBasicAuth() .appendHeader( Headers.ACCEPT, APPLICATION_JSON, @@ -221,6 +242,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr override suspend fun blobStream(repository: Repository, digest: Digest): Result { val deserializable = SourceDeserializer() return fuelManager.get("/v2/$repository/blobs/$digest") + .addBasicAuth() .appendHeader( Headers.ACCEPT, APPLICATION_JSON, @@ -237,6 +259,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr override suspend fun initiateUpload(repository: Repository): Result = fuelManager.post("/v2/$repository/blobs/uploads/") + .addBasicAuth() .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .mapToUploadSession() @@ -245,6 +268,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr val parameters = listOf("digest" to digest) return fuelManager.put(session.location, parameters) + .addBasicAuth() .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .mapToDigest() @@ -256,6 +280,7 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr startRange: Long, endRange: Long, ): Result = fuelManager.patch(session.location) + .addBasicAuth() .appendHeader(Headers.CONTENT_LENGTH, buffer.size) .appendHeader("Content-Range", "$startRange-$endRange") .appendHeader(Headers.CONTENT_TYPE, APPLICATION_OCTET_STREAM) @@ -269,22 +294,30 @@ internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, cr path: Path, size: Long, digest: Digest, - ): Result = - fuelManager.put(session.location, listOf("digest" to digest)) - .appendHeader(Headers.CONTENT_TYPE, APPLICATION_OCTET_STREAM) - .body(path::inputStream, { size }, repeatable = true) - .awaitResponseResult(EmptyDeserializer) - .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } - .mapToDigest() + ): Result = fuelManager.put(session.location, listOf("digest" to digest)) + .addBasicAuth() + .appendHeader(Headers.CONTENT_TYPE, APPLICATION_OCTET_STREAM) + .body(path::inputStream, { size }) + .awaitResponseResult(EmptyDeserializer) + .let { responseResult -> + handler.retryOnUnauthorized( + responseResult, + EmptyDeserializer, + RetryMode.Stream(path::inputStream, size), + ) + } + .mapToDigest() override suspend fun uploadStatus(session: UploadSession): Result, FuelError> = fuelManager.get(session.location) + .addBasicAuth() .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .mapToRange() override suspend fun cancelBlobUpload(session: UploadSession): Result<*, FuelError> = fuelManager.delete(session.location) + .addBasicAuth() .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .third diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ResponseRetryWithAuthentication.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ResponseRetryWithAuthentication.kt index 682bcf35..bc77f067 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ResponseRetryWithAuthentication.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ResponseRetryWithAuthentication.kt @@ -14,6 +14,18 @@ import de.cmdjulian.kirc.client.RegistryCredentials import de.cmdjulian.kirc.utils.CaseInsensitiveMap import im.toss.http.parser.HttpAuthCredentials import io.goodforgod.graalvm.hint.annotation.ReflectionHint +import java.io.InputStream + +/** + * Defines body source for the retry method: + * + * - [Default]: copies body from original request + * - [Stream]: adds InputStream as body from provided lambda + */ +internal sealed class RetryMode { + data object Default : RetryMode() + data class Stream(val stream: () -> InputStream, val size: Long) : RetryMode() +} internal class ResponseRetryWithAuthentication( private val credentials: RegistryCredentials?, @@ -22,12 +34,17 @@ internal class ResponseRetryWithAuthentication( suspend fun retryOnUnauthorized( responseResult: ResponseResultOf, deserializer: Deserializable, + mode: RetryMode = RetryMode.Default, ): ResponseResultOf { val (request, response, _) = responseResult val headers = CaseInsensitiveMap(response.headers) if (response.statusCode == 401 && "www-authenticate" in headers) { val retryableRequest = retryRequest(headers["www-authenticate"]?.first(), request) + when (mode) { + is RetryMode.Default -> Unit + is RetryMode.Stream -> retryableRequest?.body(mode.stream, { mode.size }) + } retryableRequest?.let { return it.awaitResponseResult(deserializer) } } From e6e79bc5497cd97ed8732108c0a82e26137d7d75 Mon Sep 17 00:00:00 2001 From: Mirko Hirsch Date: Mon, 15 Dec 2025 16:36:26 +0100 Subject: [PATCH 5/7] refactor ktlint --- .../src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt | 3 +-- .../kirc/client/ReactiveContainerImageRegistryClient.kt | 1 - .../kirc/client/SuspendingContainerImageRegistryClient.kt | 1 - .../src/main/kotlin/de/cmdjulian/kirc/utils/UtilityFuns.kt | 2 +- .../src/test/kotlin/de/cmdjulian/kirc/BlockingRegistryTest.kt | 1 - .../test/kotlin/de/cmdjulian/kirc/DockerRegistryCliHelper.kt | 1 - 6 files changed, 2 insertions(+), 7 deletions(-) diff --git a/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt b/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt index 05e13157..27906639 100644 --- a/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt +++ b/kirc-core/src/main/kotlin/de/cmdjulian/kirc/client/UploadMode.kt @@ -13,6 +13,5 @@ sealed class UploadMode { @Deprecated("Use Chunked and specify the chunk size") data object Compatibility : UploadMode() - data class Chunked(val chunkSize: Long) : UploadMode() -} \ No newline at end of file +} diff --git a/kirc-reactive/src/main/kotlin/de/cmdjulian/kirc/client/ReactiveContainerImageRegistryClient.kt b/kirc-reactive/src/main/kotlin/de/cmdjulian/kirc/client/ReactiveContainerImageRegistryClient.kt index cc9b4686..526f303e 100644 --- a/kirc-reactive/src/main/kotlin/de/cmdjulian/kirc/client/ReactiveContainerImageRegistryClient.kt +++ b/kirc-reactive/src/main/kotlin/de/cmdjulian/kirc/client/ReactiveContainerImageRegistryClient.kt @@ -4,7 +4,6 @@ import de.cmdjulian.kirc.image.Digest import de.cmdjulian.kirc.image.Reference import de.cmdjulian.kirc.image.Repository import de.cmdjulian.kirc.image.Tag - import de.cmdjulian.kirc.spec.image.ImageConfig import de.cmdjulian.kirc.spec.manifest.Manifest import de.cmdjulian.kirc.spec.manifest.ManifestList diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/client/SuspendingContainerImageRegistryClient.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/client/SuspendingContainerImageRegistryClient.kt index 80ba5322..9906964e 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/client/SuspendingContainerImageRegistryClient.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/client/SuspendingContainerImageRegistryClient.kt @@ -4,7 +4,6 @@ import de.cmdjulian.kirc.image.Digest import de.cmdjulian.kirc.image.Reference import de.cmdjulian.kirc.image.Repository import de.cmdjulian.kirc.image.Tag - import de.cmdjulian.kirc.impl.response.ResultSource import de.cmdjulian.kirc.impl.response.UploadSession import de.cmdjulian.kirc.spec.UploadBlobPath diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/utils/UtilityFuns.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/utils/UtilityFuns.kt index 35b94a76..58d005c9 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/utils/UtilityFuns.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/utils/UtilityFuns.kt @@ -1,7 +1,7 @@ package de.cmdjulian.kirc.utils import kotlin.io.path.pathString -import kotlinx.io.files.Path as KotlinPath import java.nio.file.Path as JavaPath +import kotlinx.io.files.Path as KotlinPath internal fun JavaPath.toKotlinPath(): KotlinPath = KotlinPath(pathString) diff --git a/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/BlockingRegistryTest.kt b/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/BlockingRegistryTest.kt index b62aabf5..1eb4ea10 100644 --- a/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/BlockingRegistryTest.kt +++ b/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/BlockingRegistryTest.kt @@ -7,7 +7,6 @@ import de.cmdjulian.kirc.client.UploadMode import de.cmdjulian.kirc.image.Digest import de.cmdjulian.kirc.image.Repository import de.cmdjulian.kirc.image.Tag - import de.cmdjulian.kirc.spec.manifest.ManifestList import io.kotest.assertions.throwables.shouldNotThrowAny import io.kotest.matchers.booleans.shouldBeTrue diff --git a/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/DockerRegistryCliHelper.kt b/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/DockerRegistryCliHelper.kt index 39d0ec12..bde1fbd2 100644 --- a/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/DockerRegistryCliHelper.kt +++ b/kirc-suspending/src/test/kotlin/de/cmdjulian/kirc/DockerRegistryCliHelper.kt @@ -7,7 +7,6 @@ import de.cmdjulian.kirc.client.UploadMode import de.cmdjulian.kirc.image.Digest import de.cmdjulian.kirc.image.Reference import de.cmdjulian.kirc.image.Repository - import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.runBlocking import kotlinx.io.asInputStream From a744ef6e0c247c573ae449b11c945b2b366e686b Mon Sep 17 00:00:00 2001 From: Mirko Hirsch Date: Tue, 16 Dec 2025 16:22:42 +0100 Subject: [PATCH 6/7] refactor --- .../kirc/impl/ContainerRegistryApiImpl.kt | 30 +------------------ 1 file changed, 1 insertion(+), 29 deletions(-) diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt index 86df81ba..ca715c84 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/ContainerRegistryApiImpl.kt @@ -4,11 +4,9 @@ import com.github.kittinunf.fuel.core.FuelError import com.github.kittinunf.fuel.core.FuelManager import com.github.kittinunf.fuel.core.Headers import com.github.kittinunf.fuel.core.Parameters -import com.github.kittinunf.fuel.core.Request import com.github.kittinunf.fuel.core.awaitResponseResult import com.github.kittinunf.fuel.core.deserializers.ByteArrayDeserializer import com.github.kittinunf.fuel.core.deserializers.EmptyDeserializer -import com.github.kittinunf.fuel.core.extensions.authentication import com.github.kittinunf.result.Result import com.github.kittinunf.result.flatMap import com.github.kittinunf.result.map @@ -45,22 +43,14 @@ import kotlin.io.path.inputStream private const val APPLICATION_JSON = "application/json" private const val APPLICATION_OCTET_STREAM = "application/octet-stream" -internal class ContainerRegistryApiImpl( - private val fuelManager: FuelManager, - private val credentials: RegistryCredentials?, -) : +internal class ContainerRegistryApiImpl(private val fuelManager: FuelManager, credentials: RegistryCredentials?) : ContainerRegistryApi { private val handler = ResponseRetryWithAuthentication(credentials, fuelManager) - private fun Request.addBasicAuth() = apply { - credentials?.let { credentials -> authentication().basic(credentials.username, credentials.password) } - } - // Status override suspend fun ping(): Result<*, FuelError> = fuelManager.get("/v2/") - .addBasicAuth() .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .third @@ -73,7 +63,6 @@ internal class ContainerRegistryApiImpl( val deserializable = jacksonDeserializer() return fuelManager.get("/v2/_catalog", parameter) - .addBasicAuth() .appendHeader(Headers.ACCEPT, APPLICATION_JSON) .awaitResponseResult(deserializable) .let { responseResult -> handler.retryOnUnauthorized(responseResult, deserializable) } @@ -88,7 +77,6 @@ internal class ContainerRegistryApiImpl( val deserializable = jacksonDeserializer() return fuelManager.get("/v2/$repository/tags/list", parameter) - .addBasicAuth() .appendHeader(Headers.ACCEPT, APPLICATION_JSON) .awaitResponseResult(deserializable) .let { responseResult -> handler.retryOnUnauthorized(responseResult, deserializable) } @@ -98,7 +86,6 @@ internal class ContainerRegistryApiImpl( override suspend fun digest(repository: Repository, reference: Reference) = when (reference) { is Digest -> Result.success(reference) is Tag -> fuelManager.head("/v2/$repository/manifests/$reference") - .addBasicAuth() .appendHeader( Headers.ACCEPT, APPLICATION_JSON, @@ -119,7 +106,6 @@ internal class ContainerRegistryApiImpl( reference: Reference, accept: String, ): Result<*, FuelError> = fuelManager.head("/v2/$repository/manifests/$reference") - .addBasicAuth() .appendHeader(Headers.ACCEPT, accept) .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } @@ -128,7 +114,6 @@ internal class ContainerRegistryApiImpl( override suspend fun manifests(repository: Repository, reference: Reference): Result { val deserializable = jacksonDeserializer() return fuelManager.get("/v2/$repository/manifests/$reference") - .addBasicAuth() .appendHeader( Headers.ACCEPT, APPLICATION_JSON, @@ -145,7 +130,6 @@ internal class ContainerRegistryApiImpl( override suspend fun manifestStream(repository: Repository, reference: Reference): Result { val deserializable = SourceDeserializer() return fuelManager.get("/v2/$repository/manifests/$reference") - .addBasicAuth() .appendHeader( Headers.ACCEPT, APPLICATION_JSON, @@ -162,7 +146,6 @@ internal class ContainerRegistryApiImpl( override suspend fun manifest(repository: Repository, reference: Reference): Result { val deserializable = jacksonDeserializer() return fuelManager.get("/v2/$repository/manifests/$reference") - .addBasicAuth() .appendHeader(Headers.ACCEPT, APPLICATION_JSON, OciManifestV1.MediaType, DockerManifestV2.MediaType) .awaitResponseResult(deserializable) .let { responseResult -> handler.retryOnUnauthorized(responseResult, deserializable) } @@ -185,7 +168,6 @@ internal class ContainerRegistryApiImpl( is OciManifestListV1 -> OciManifestListV1.MediaType } return fuelManager.put("/v2/$repository/manifests/$urlReference") - .addBasicAuth() .appendHeader(Headers.CONTENT_TYPE, contentType) .body(JsonMapper.writeValueAsString(manifest)) .awaitResponseResult(EmptyDeserializer) @@ -196,7 +178,6 @@ internal class ContainerRegistryApiImpl( override suspend fun deleteManifest(repository: Repository, reference: Reference): Result { return digest(repository, reference).flatMap { digest -> fuelManager.delete("/v2/$repository/manifests/$digest") - .addBasicAuth() .appendHeader( Headers.ACCEPT, APPLICATION_JSON, @@ -216,7 +197,6 @@ internal class ContainerRegistryApiImpl( override suspend fun existsBlob(repository: Repository, digest: Digest): Result<*, FuelError> = fuelManager.head("/v2/$repository/blobs/$digest") - .addBasicAuth() .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .third @@ -224,7 +204,6 @@ internal class ContainerRegistryApiImpl( override suspend fun blob(repository: Repository, digest: Digest): Result { val deserializable = ByteArrayDeserializer() return fuelManager.get("/v2/$repository/blobs/$digest") - .addBasicAuth() .appendHeader( Headers.ACCEPT, APPLICATION_JSON, @@ -242,7 +221,6 @@ internal class ContainerRegistryApiImpl( override suspend fun blobStream(repository: Repository, digest: Digest): Result { val deserializable = SourceDeserializer() return fuelManager.get("/v2/$repository/blobs/$digest") - .addBasicAuth() .appendHeader( Headers.ACCEPT, APPLICATION_JSON, @@ -259,7 +237,6 @@ internal class ContainerRegistryApiImpl( override suspend fun initiateUpload(repository: Repository): Result = fuelManager.post("/v2/$repository/blobs/uploads/") - .addBasicAuth() .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .mapToUploadSession() @@ -268,7 +245,6 @@ internal class ContainerRegistryApiImpl( val parameters = listOf("digest" to digest) return fuelManager.put(session.location, parameters) - .addBasicAuth() .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .mapToDigest() @@ -280,7 +256,6 @@ internal class ContainerRegistryApiImpl( startRange: Long, endRange: Long, ): Result = fuelManager.patch(session.location) - .addBasicAuth() .appendHeader(Headers.CONTENT_LENGTH, buffer.size) .appendHeader("Content-Range", "$startRange-$endRange") .appendHeader(Headers.CONTENT_TYPE, APPLICATION_OCTET_STREAM) @@ -295,7 +270,6 @@ internal class ContainerRegistryApiImpl( size: Long, digest: Digest, ): Result = fuelManager.put(session.location, listOf("digest" to digest)) - .addBasicAuth() .appendHeader(Headers.CONTENT_TYPE, APPLICATION_OCTET_STREAM) .body(path::inputStream, { size }) .awaitResponseResult(EmptyDeserializer) @@ -310,14 +284,12 @@ internal class ContainerRegistryApiImpl( override suspend fun uploadStatus(session: UploadSession): Result, FuelError> = fuelManager.get(session.location) - .addBasicAuth() .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .mapToRange() override suspend fun cancelBlobUpload(session: UploadSession): Result<*, FuelError> = fuelManager.delete(session.location) - .addBasicAuth() .awaitResponseResult(EmptyDeserializer) .let { responseResult -> handler.retryOnUnauthorized(responseResult, EmptyDeserializer) } .third From 89b147b919549cc554dd45a2f20fa1aa51843332 Mon Sep 17 00:00:00 2001 From: Mirko Hirsch Date: Tue, 16 Dec 2025 17:10:22 +0100 Subject: [PATCH 7/7] refactor --- .../SuspendingContainerImageRegistryClientImpl.kt | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/SuspendingContainerImageRegistryClientImpl.kt b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/SuspendingContainerImageRegistryClientImpl.kt index 897c0c84..9e9c6f1b 100644 --- a/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/SuspendingContainerImageRegistryClientImpl.kt +++ b/kirc-suspending/src/main/kotlin/de/cmdjulian/kirc/impl/SuspendingContainerImageRegistryClientImpl.kt @@ -143,11 +143,13 @@ internal class SuspendingContainerImageRegistryClientImpl(private val api: Conta } catch (_: EOFException) { // expected behavior when EOF is reached } finally { - val bytesRead = buffer.size - val endRange = startRange + bytesRead - 1 - currentSession = api.uploadBlobChunked(currentSession, buffer, startRange, endRange) - .getOrElse { throw it.toRegistryClientError() } - startRange = endRange + 1 + if (buffer.size > 0) { + val bytesRead = buffer.size + val endRange = startRange + bytesRead - 1 + currentSession = api.uploadBlobChunked(currentSession, buffer, startRange, endRange) + .getOrElse { throw it.toRegistryClientError() } + startRange = endRange + 1 + } } } currentSession