diff --git a/ROADMAP.md b/ROADMAP.md
index 6e533c97..573a576d 100644
--- a/ROADMAP.md
+++ b/ROADMAP.md
@@ -217,6 +217,7 @@ Order = user-visible value first; each ships with backend + UI + tests.
- [x] Share `hide_readme` option (upstream #2729 item 6) — `ShareProps.HideReadMe` (only meaningful with `ShowReadMe`); share navigator filters `README.md`/`README.txt` (case-insensitive) from listings while direct-path resolution stays open for the readme viewer; `detectReadMe` URI fallback now probes unconditionally; owner-only `hide_readme` on share responses; Share dialog nested checkbox
- [x] 2FA recovery codes (upstream #2729 item 3) — `user.two_factor_backup_codes` sensitive JSON of `salt:sha256` digests; `PUT /user/setting/2fa/backup` regenerates 10 one-time codes behind a valid TOTP (rate-limited 5/h); `Verify2FA` falls back to single-use code consumption on TOTP failure; codes invalidated on secret rotation/disable; login phase gains a recovery-code input mode; security settings show remaining count + regenerate dialog
- [x] Public share directory (upstream #2729 items 4+5) — `share.listed_publicly` opt-in column gated by new `GroupPermissionSharePublicList` group bit; rejected on password-protected shares at the service layer and normalized off at creation; `GET /share/listed` anonymous endpoint (rate-limited 60/min/IP, cursor pagination) listing only non-expired passwordless listed shares with case-insensitive name search across anchor + covered files; `listed_publicly` owner-visible in share responses; `/discover` page (anonymous-visible nav item + sign-in link), admin group Share section switch, en+zh locales
+- [x] Decompression-bomb guards (meta #2 item 10) — `DecompressSize` now bounds cumulative extracted output (its documented "total file size" intent), not just compressed input: `checkExtractGuards` aborts at the limit and at a 100k-entry cap (`maxExtractEntries`, bounds dir-creation bombs), each entry stream wrapped in `cappedFile` so understated size headers cannot overrun; slave path receives the limit via `SlaveExtractArchiveTaskState.ExtractLimit`; all failures carry `queue.CriticalErr` (no retry of the same bomb); resume-safe via cursor-skip size accounting
- [x] Storage policy total capacity (upstream #2178 item 1) — `PolicySetting.MaxTotalSize` caps cumulative entity bytes per policy; enforced in `PrepareUpload`, batch upload validation, and `copyFiles` (baseline usage + per-batch accumulation since tx writes are invisible to the usage query); canonical `ErrInsufficientCapacity` preserved so `errors.Is` quota handling still matches; admin policy editor gains a Max total capacity SizeInput, en+zh locales
## 6. Phase D — desktop, all platforms
@@ -241,7 +242,7 @@ Lives in `android/` in this repo. Kotlin + Jetpack Compose, Material 3.
- **API**: `api/v4` REST + OAuth token (entities exist: `oauthclient`, `oauthgrant`) — same surface the desktop `cloudreve-api` crate documents; port its models as the spec
- **Core features**: browse/download/upload files, share links, full-text search (done — query bar + offset pagination + parent-path/snippet rows), camera-upload (done — periodic WorkManager MediaStore sync, Wi-Fi-only constraint, ID dedup, settings dialog), offline-favorite files (done — "Keep offline" downloads to filesDir/offline, DataStore registry, star dialog with open/refresh/remove), local sync folder via SAF/WorkManager
-- **System integration** (the "native, complete" ask): share-sheet target (upload to Cloudreve from any app), DocumentsProvider (Cloudreve in Files app), quick-share tile, notifications on share/task events
+- **System integration** (the "native, complete" ask): share-sheet target (done — SEND/SEND_MULTIPLE → UploadWorker), DocumentsProvider (done — Files-app browse/open/thumb/rename/delete/search), quick-share tile, notifications on share/task events
- **Auth**: webview OAuth flow → token; later passkey if backend exposes
- **WebDAV bridge**: `/dav` works as fallback file access until SDK matures
- Non-goals: iOS, tablet-first layouts (works, not optimized)
diff --git a/android/README.md b/android/README.md
index 804ed214..f446fad2 100644
--- a/android/README.md
+++ b/android/README.md
@@ -38,8 +38,11 @@ CI runs `assembleDebug` on every PR.
`filesDir/offline/` and registers the entry (DataStore JSON);
star icon in the top bar opens the offline list — open via
FileProvider, re-download to refresh, remove to delete
+- DocumentsProvider: the whole tree appears in the system Files app /
+ SAF pickers — browse, open (cached download), thumbnails, rename,
+ delete, and search all proxy to `api/v4`
## Planned next
-- DocumentsProvider
+- Quick-share tile, task/share notifications, local sync folder
- No iOS. Ever.
diff --git a/android/app/src/main/AndroidManifest.xml b/android/app/src/main/AndroidManifest.xml
index 7ac9fbc7..23b3e9ac 100644
--- a/android/app/src/main/AndroidManifest.xml
+++ b/android/app/src/main/AndroidManifest.xml
@@ -47,5 +47,16 @@
android:name="android.support.FILE_PROVIDER_PATHS"
android:resource="@xml/file_paths" />
+
+
+
+
+
+
diff --git a/android/app/src/main/java/org/cloudreve/android/provider/CloudreveDocumentsProvider.kt b/android/app/src/main/java/org/cloudreve/android/provider/CloudreveDocumentsProvider.kt
new file mode 100644
index 00000000..bfad6c8f
--- /dev/null
+++ b/android/app/src/main/java/org/cloudreve/android/provider/CloudreveDocumentsProvider.kt
@@ -0,0 +1,247 @@
+package org.cloudreve.android.provider
+
+import android.content.Context
+import android.database.Cursor
+import android.database.MatrixCursor
+import android.graphics.Point
+import android.os.CancellationSignal
+import android.os.ParcelFileDescriptor
+import android.provider.DocumentsContract
+import android.provider.DocumentsContract.Document
+import android.provider.DocumentsContract.Root
+import android.provider.DocumentsProvider
+import kotlinx.coroutines.Dispatchers
+import kotlinx.coroutines.runBlocking
+import org.cloudreve.android.CloudreveApp
+import org.cloudreve.android.R
+import org.cloudreve.android.api.FileObject
+import org.cloudreve.android.data.FileRepository
+import org.cloudreve.android.util.CrUri
+import java.io.File
+import java.io.FileNotFoundException
+import java.security.MessageDigest
+
+/**
+ * Exposes the Cloudreve filesystem to the system file picker / Files app.
+ * Document IDs are the entries' `cloudreve://` URIs, matching the desktop
+ * client's addressing scheme.
+ *
+ * Provider methods run on binder threads, so the suspend repository calls are
+ * bridged with runBlocking. All network errors surface as
+ * FileNotFoundException, which the Documents UI renders as "unavailable".
+ */
+class CloudreveDocumentsProvider : DocumentsProvider() {
+
+ private val repo: FileRepository
+ get() = (context!!.applicationContext as CloudreveApp).fileRepository
+
+ private val ioScope = Dispatchers.IO
+
+ override fun onCreate(): Boolean = true
+
+ override fun queryRoots(projection: Array?): Cursor {
+ val cursor = MatrixCursor(projection ?: DEFAULT_ROOT_COLUMNS)
+ cursor.newRow().apply {
+ add(Root.COLUMN_ROOT_ID, ROOT_ID)
+ add(Root.COLUMN_DOCUMENT_ID, ROOT_ID)
+ add(Root.COLUMN_TITLE, "Cloudreve")
+ add(Root.COLUMN_SUMMARY, rootSummary())
+ add(Root.COLUMN_FLAGS, Root.FLAG_SUPPORTS_SEARCH or Root.FLAG_SUPPORTS_IS_CHILD)
+ add(Root.COLUMN_ICON, R.mipmap.ic_launcher)
+ add(Root.COLUMN_MIME_TYPES, "*/*")
+ }
+ return cursor
+ }
+
+ private fun rootSummary(): String = runCatching {
+ val base = runBlocking(ioScope) {
+ (context!!.applicationContext as CloudreveApp).apiClient.serverBase()
+ }
+ java.net.URI(base).host
+ }.getOrNull() ?: ""
+
+ override fun isChildDocument(parentDocumentId: String, documentId: String): Boolean =
+ documentId != parentDocumentId && documentId.startsWith("${parentDocumentId.trimEnd('/')}/")
+
+ override fun queryDocument(documentId: String, projection: Array?): Cursor {
+ val cursor = MatrixCursor(projection ?: DEFAULT_DOC_COLUMNS)
+ if (documentId == ROOT_ID) {
+ cursor.newRow().apply {
+ add(Document.COLUMN_DOCUMENT_ID, ROOT_ID)
+ add(Document.COLUMN_DISPLAY_NAME, "My files")
+ add(Document.COLUMN_MIME_TYPE, Document.MIME_TYPE_DIR)
+ add(Document.COLUMN_FLAGS, 0)
+ }
+ return cursor
+ }
+ val obj = findObject(documentId)
+ cursor.addFileRow(obj)
+ return cursor
+ }
+
+ override fun queryChildDocuments(
+ parentDocumentId: String,
+ projection: Array?,
+ sortOrder: String?,
+ ): Cursor {
+ val cursor = MatrixCursor(projection ?: DEFAULT_DOC_COLUMNS)
+ runBlocking(ioScope) {
+ var token: String? = null
+ do {
+ val page = repo.list(parentDocumentId, token)
+ page.files.forEach { cursor.addFileRow(it) }
+ token = page.pagination.nextToken
+ } while (token != null)
+ }
+ return cursor
+ }
+
+ override fun querySearchDocuments(
+ rootId: String,
+ query: String,
+ projection: Array?,
+ ): Cursor {
+ val cursor = MatrixCursor(projection ?: DEFAULT_DOC_COLUMNS)
+ runBlocking(ioScope) {
+ repo.search(query).hits.forEach { cursor.addFileRow(it.file) }
+ }
+ return cursor
+ }
+
+ override fun openDocument(
+ documentId: String,
+ mode: String,
+ signal: CancellationSignal?,
+ ): ParcelFileDescriptor {
+ if (mode.contains('w')) {
+ throw FileNotFoundException("Cloudreve documents are read-only here")
+ }
+ val local = fetchToCache(documentId)
+ return ParcelFileDescriptor.open(local, ParcelFileDescriptor.MODE_READ_ONLY)
+ }
+
+ override fun openDocumentThumbnail(
+ documentId: String,
+ sizeHint: Point?,
+ signal: CancellationSignal?,
+ ): android.content.res.AssetFileDescriptor {
+ val local = runBlocking(ioScope) {
+ val url = repo.thumbUrl(documentId)
+ ?: throw FileNotFoundException("No thumbnail")
+ downloadTo(url, cacheFile("thumb", documentId))
+ }
+ val pfd = ParcelFileDescriptor.open(local, ParcelFileDescriptor.MODE_READ_ONLY)
+ return android.content.res.AssetFileDescriptor(pfd, 0, local.length())
+ }
+
+ override fun deleteDocument(documentId: String) {
+ runBlocking(ioScope) { repo.delete(listOf(documentId)) }
+ notifyChange(documentId)
+ }
+
+ override fun renameDocument(documentId: String, displayName: String): String {
+ runBlocking(ioScope) { repo.rename(documentId, displayName) }
+ val newId = CrUri.join(CrUri.parent(documentId), displayName)
+ notifyChange(documentId)
+ return newId
+ }
+
+ private fun notifyChange(documentId: String) {
+ context!!.contentResolver.notifyChange(
+ DocumentsContract.buildDocumentUri(authority(), documentId),
+ null,
+ )
+ }
+
+ private fun authority(): String = "${context!!.packageName}.documents"
+
+ /** Looks up a file by listing its parent directory. */
+ private fun findObject(documentId: String): FileObject = runBlocking(ioScope) {
+ val parent = CrUri.parent(documentId)
+ var token: String? = null
+ do {
+ val page = repo.list(parent, token)
+ page.files.firstOrNull { it.path == documentId }?.let { return@runBlocking it }
+ token = page.pagination.nextToken
+ } while (token != null)
+ throw FileNotFoundException("Not found: $documentId")
+ }
+
+ private fun fetchToCache(documentId: String): File = runBlocking(ioScope) {
+ val url = repo.downloadUrl(documentId)
+ downloadTo(url, cacheFile("doc", documentId))
+ }
+
+ private suspend fun downloadTo(url: String, dest: File): File {
+ val resp = repo.download(url)
+ if (!resp.isSuccessful) {
+ resp.body()?.close()
+ throw FileNotFoundException("Download failed: HTTP ${resp.code()}")
+ }
+ dest.parentFile?.mkdirs()
+ resp.body()!!.byteStream().use { input ->
+ dest.outputStream().use { output -> input.copyTo(output) }
+ }
+ return dest
+ }
+
+ private fun cacheFile(kind: String, documentId: String): File {
+ val hash = MessageDigest.getInstance("SHA-256")
+ .digest(documentId.toByteArray())
+ .joinToString("") { "%02x".format(it) }
+ .take(16)
+ val ext = CrUri.fileName(documentId).substringAfterLast('.', "")
+ val suffix = if (ext.isEmpty()) "" else ".$ext"
+ return File(context!!.cacheDir, "docprovider/$kind-$hash$suffix")
+ }
+
+ private fun MatrixCursor.addFileRow(file: FileObject) {
+ val isDir = file.isFolder
+ val flags = Document.FLAG_SUPPORTS_DELETE or
+ Document.FLAG_SUPPORTS_RENAME or
+ (if (!isDir && file.thumbnail == true) Document.FLAG_SUPPORTS_THUMBNAIL else 0)
+ newRow().apply {
+ add(Document.COLUMN_DOCUMENT_ID, file.path)
+ add(Document.COLUMN_DISPLAY_NAME, file.name)
+ add(
+ Document.COLUMN_MIME_TYPE,
+ if (isDir) Document.MIME_TYPE_DIR
+ else java.net.URLConnection.guessContentTypeFromName(file.name)
+ ?: "application/octet-stream",
+ )
+ add(Document.COLUMN_SIZE, if (isDir) null else file.size)
+ add(Document.COLUMN_LAST_MODIFIED, parseInstant(file.updatedAt))
+ add(Document.COLUMN_FLAGS, flags)
+ }
+ }
+
+ private fun parseInstant(value: String): Long? {
+ if (value.isBlank()) return null
+ return runCatching { java.time.Instant.parse(value).toEpochMilli() }
+ .recoverCatching { java.time.OffsetDateTime.parse(value).toInstant().toEpochMilli() }
+ .getOrNull()
+ }
+
+ companion object {
+ private const val ROOT_ID = CrUri.MY_PREFIX
+
+ private val DEFAULT_ROOT_COLUMNS = arrayOf(
+ Root.COLUMN_ROOT_ID,
+ Root.COLUMN_DOCUMENT_ID,
+ Root.COLUMN_TITLE,
+ Root.COLUMN_SUMMARY,
+ Root.COLUMN_FLAGS,
+ Root.COLUMN_ICON,
+ Root.COLUMN_MIME_TYPES,
+ )
+
+ private val DEFAULT_DOC_COLUMNS = arrayOf(
+ Document.COLUMN_DOCUMENT_ID,
+ Document.COLUMN_DISPLAY_NAME,
+ Document.COLUMN_MIME_TYPE,
+ Document.COLUMN_SIZE,
+ Document.COLUMN_LAST_MODIFIED,
+ Document.COLUMN_FLAGS,
+ )
+ }
+}
diff --git a/pkg/filemanager/workflows/extract.go b/pkg/filemanager/workflows/extract.go
index 840bd4a8..7dd1e8a0 100644
--- a/pkg/filemanager/workflows/extract.go
+++ b/pkg/filemanager/workflows/extract.go
@@ -5,6 +5,7 @@ import (
"encoding/json"
"fmt"
"io"
+ iofs "io/fs"
"os"
"path"
"path/filepath"
@@ -65,6 +66,12 @@ const (
ProgressTypeExtractSize = "extract_size"
ProgressTypeDownload = "download"
+ // maxExtractEntries bounds the number of archive entries a single
+ // extraction task will process. Archives with more entries than this
+ // abort — a decompression bomb with millions of tiny entries would
+ // otherwise flood the file table.
+ maxExtractEntries int64 = 100_000
+
SummaryKeySrc = "src"
SummaryKeySrcPhysical = "src_physical"
SummaryKeyDst = "dst"
@@ -212,15 +219,16 @@ func (m *ExtractArchiveTask) createSlaveExtractTask(ctx context.Context, dep dep
}
payload := &SlaveExtractArchiveTaskState{
- FileName: archiveFile.DisplayName(),
- Entity: entityModel,
- Policy: policy,
- Encoding: m.state.Encoding,
- Dst: m.state.Dst,
- UserID: user.ID,
- Password: m.state.Password,
- FileMask: m.state.FileMask,
- Volumes: m.resolveVolumeEntities(ctx, fm, uri.DirUri(), archiveFile.DisplayName(), entityModel),
+ FileName: archiveFile.DisplayName(),
+ Entity: entityModel,
+ Policy: policy,
+ Encoding: m.state.Encoding,
+ Dst: m.state.Dst,
+ UserID: user.ID,
+ ExtractLimit: user.Edges.Group.Settings.DecompressSize,
+ Password: m.state.Password,
+ FileMask: m.state.FileMask,
+ Volumes: m.resolveVolumeEntities(ctx, fm, uri.DirUri(), archiveFile.DisplayName(), entityModel),
}
payloadStr, err := json.Marshal(payload)
@@ -424,6 +432,13 @@ func (m *ExtractArchiveTask) masterExtractArchive(ctx context.Context, dep depen
return nil
}
+ // Decompression-bomb guard: abort when cumulative output or entry
+ // count exceeds the group's bounds.
+ sizeLimit := user.Edges.Group.Settings.DecompressSize
+ if err := checkExtractGuards(m.progress, sizeLimit); err != nil {
+ return err
+ }
+
if f.FileInfo.IsDir() {
_, err := fm.Create(ctx, savePath, types.FileTypeFolder)
if err != nil {
@@ -441,6 +456,13 @@ func (m *ExtractArchiveTask) masterExtractArchive(ctx context.Context, dep depen
return nil
}
+ if sizeLimit > 0 {
+ // Declared entry sizes are advisory; cap the stream at the
+ // remaining budget so understated sizes cannot overrun.
+ remaining := sizeLimit - atomic.LoadInt64(&m.progress[ProgressTypeExtractSize].Current)
+ fileStream = &cappedFile{File: fileStream, remaining: remaining}
+ }
+
fileData := &fs.UploadRequest{
Props: &fs.UploadProps{
Uri: savePath,
@@ -660,6 +682,7 @@ type (
Encoding string `json:"encoding,omitempty"`
Dst string `json:"dst,omitempty"`
UserID int `json:"user_id"`
+ ExtractLimit int64 `json:"extract_limit,omitempty"`
TempPath string `json:"temp_path,omitempty"`
TempZipFilePath string `json:"temp_zip_file_path,omitempty"`
ProcessedCursor string `json:"processed_cursor,omitempty"`
@@ -870,6 +893,10 @@ func (m *SlaveExtractArchiveTask) Do(ctx context.Context) (task.Status, error) {
return nil
}
+ if err := checkExtractGuards(m.progress, m.state.ExtractLimit); err != nil {
+ return err
+ }
+
if f.FileInfo.IsDir() {
_, err := fm.Create(ctx, savePath, types.FileTypeFolder, fs.WithNode(m.node), fs.WithStatelessUserID(m.state.UserID))
if err != nil {
@@ -887,6 +914,11 @@ func (m *SlaveExtractArchiveTask) Do(ctx context.Context) (task.Status, error) {
return nil
}
+ if m.state.ExtractLimit > 0 {
+ remaining := m.state.ExtractLimit - atomic.LoadInt64(&m.progress[ProgressTypeExtractSize].Current)
+ fileStream = &cappedFile{File: fileStream, remaining: remaining}
+ }
+
fileData := &fs.UploadRequest{
Props: &fs.UploadProps{
Uri: savePath,
@@ -947,3 +979,55 @@ func isFileInMask(path string, mask []string) bool {
return false
}
+
+// errExtractSizeLimit aborts extraction when cumulative decompressed output
+// exceeds the group's DecompressSize bound. It carries CriticalErr so retries
+// do not reprocess the same bomb.
+var errExtractSizeLimit = fmt.Errorf("extracted size exceeds the decompress limit: %w", queue.CriticalErr)
+
+// checkExtractGuards enforces the decompression-bomb bounds before an archive
+// entry is written: cumulative output size (declared, from the shared progress
+// counter) and total entry count. A returned error aborts the task as a
+// critical failure — no retry will change the outcome.
+func checkExtractGuards(progress queue.Progresses, sizeLimit int64) error {
+ if sizeLimit > 0 {
+ current := atomic.LoadInt64(&progress[ProgressTypeExtractSize].Current)
+ if current >= sizeLimit {
+ return fmt.Errorf("%w (%d >= %d)", errExtractSizeLimit, current, sizeLimit)
+ }
+ }
+ if count := atomic.LoadInt64(&progress[ProgressTypeExtractCount].Current); count >= maxExtractEntries {
+ return fmt.Errorf("archive exceeds the entry limit %d: %w", maxExtractEntries, queue.CriticalErr)
+ }
+ return nil
+}
+
+// cappedFile bounds a single archive entry's stream at `remaining`
+// bytes. Declared entry sizes are advisory — a crafted archive can understate
+// them — so the stream itself is capped; reading past the budget fails the
+// upload and aborts extraction.
+type cappedFile struct {
+ iofs.File
+ remaining int64
+}
+
+func (c *cappedFile) Read(p []byte) (int, error) {
+ if len(p) == 0 {
+ return 0, nil
+ }
+ if c.remaining <= 0 {
+ // Budget exhausted: an entry ending exactly at the boundary must
+ // still see EOF, while any further data fails the upload.
+ n, err := c.File.Read(p[:1])
+ if n > 0 {
+ return 0, errExtractSizeLimit
+ }
+ return 0, err
+ }
+ if int64(len(p)) > c.remaining {
+ p = p[:c.remaining]
+ }
+ n, err := c.File.Read(p)
+ c.remaining -= int64(n)
+ return n, err
+}
diff --git a/pkg/filemanager/workflows/extract_test.go b/pkg/filemanager/workflows/extract_test.go
new file mode 100644
index 00000000..771cf94f
--- /dev/null
+++ b/pkg/filemanager/workflows/extract_test.go
@@ -0,0 +1,92 @@
+package workflows
+
+import (
+ "errors"
+ "io"
+ iofs "io/fs"
+ "strings"
+ "testing"
+
+ "github.com/cloudreve/Cloudreve/v4/pkg/queue"
+ "github.com/stretchr/testify/require"
+)
+
+func testProgress(count, size int64) queue.Progresses {
+ return queue.Progresses{
+ ProgressTypeExtractCount: &queue.Progress{Current: count},
+ ProgressTypeExtractSize: &queue.Progress{Current: size},
+ }
+}
+
+func TestCheckExtractGuards(t *testing.T) {
+ // Under all limits.
+ require.NoError(t, checkExtractGuards(testProgress(10, 100), 1000))
+
+ // At the cumulative size limit — aborts as non-retryable.
+ err := checkExtractGuards(testProgress(10, 1000), 1000)
+ require.Error(t, err)
+ require.True(t, errors.Is(err, queue.CriticalErr))
+
+ // At the entry cap — aborts as non-retryable.
+ err = checkExtractGuards(testProgress(maxExtractEntries, 10), 1000)
+ require.Error(t, err)
+ require.True(t, errors.Is(err, queue.CriticalErr))
+
+ // Zero size limit disables the size bound; entry cap still applies.
+ require.NoError(t, checkExtractGuards(testProgress(10, 1<<62), 0))
+ require.Error(t, checkExtractGuards(testProgress(maxExtractEntries, 0), 0))
+}
+
+type stubFile struct {
+ io.Reader
+}
+
+func (stubFile) Stat() (iofs.FileInfo, error) { return nil, nil }
+func (stubFile) Close() error { return nil }
+
+func TestCappedFileExactBoundary(t *testing.T) {
+ // Entry ends exactly at the budget — stream must terminate with EOF.
+ capped := &cappedFile{File: stubFile{strings.NewReader("12345")}, remaining: 5}
+ buf := make([]byte, 8)
+
+ n, err := capped.Read(buf)
+ require.NoError(t, err)
+ require.Equal(t, 5, n)
+ require.Equal(t, "12345", string(buf[:n]))
+
+ _, err = capped.Read(buf)
+ require.Equal(t, io.EOF, err)
+}
+
+func TestCappedFileBomb(t *testing.T) {
+ // Entry data beyond the budget — read fails with the critical limit error.
+ capped := &cappedFile{File: stubFile{strings.NewReader("123456")}, remaining: 5}
+ buf := make([]byte, 8)
+
+ n, err := capped.Read(buf)
+ require.NoError(t, err)
+ require.Equal(t, 5, n)
+
+ _, err = capped.Read(buf)
+ require.Error(t, err)
+ require.True(t, errors.Is(err, queue.CriticalErr))
+}
+
+func TestCappedFileShortReads(t *testing.T) {
+ // Budget caps each read; consecutive reads drain only the remaining bytes.
+ capped := &cappedFile{File: stubFile{strings.NewReader("abcdef")}, remaining: 3}
+ buf := make([]byte, 2)
+
+ var got []byte
+ var err error
+ for {
+ var n int
+ n, err = capped.Read(buf)
+ got = append(got, buf[:n]...)
+ if err != nil {
+ break
+ }
+ }
+ require.Equal(t, "abc", string(got))
+ require.True(t, errors.Is(err, queue.CriticalErr))
+}