Linux-EROFS Archive on lore.kernel.org
 help / color / mirror / Atom feed
From: Nithurshen <nithurshen.dev@gmail.com>
To: nithurshen.dev@gmail.com
Cc: hsiangkao@linux.alibaba.com, linux-erofs@lists.ozlabs.org,
	xiang@kernel.org
Subject: [PATCH 1/2 v4] fsck.erofs: add multi-threaded decompression
Date: Mon, 13 Jul 2026 12:32:49 +0530	[thread overview]
Message-ID: <20260713070249.75652-1-nithurshen.dev@gmail.com> (raw)
In-Reply-To: <20260706060525.14018-1-nithurshen.dev@gmail.com>

Currently, fsck.erofs extracts files synchronously. When decompressing
heavily packed images, the main thread spends the majority of its time
blocked on a combination of synchronous I/O syscalls and CPU-heavy
decompression routines, bottlenecking overall extraction speed.

This patch introduces a scalable, multi-threaded decompression framework
using the existing erofs_workqueue infrastructure to decouple compute
from the main thread's I/O.

To prevent massive scheduling overhead (futex contention) where worker
threads spend more CPU time waking up than actually decompressing small
clusters, this implementation introduces a batching context. Because
different compression algorithms exhibit vastly different scheduling
thresholds, the batch limit is algorithm-aware.

Signed-off-by: Nithurshen <nithurshen.dev@gmail.com>
---
 fsck/main.c              | 127 ++++++++++++---
 include/erofs/cond.h     |  31 ++++
 include/erofs/internal.h |  20 ++-
 include/erofs/lock.h     |   2 +
 lib/data.c               | 334 ++++++++++++++++++++++++++++++++++++---
 lib/global.c             |  11 +-
 6 files changed, 473 insertions(+), 52 deletions(-)
 create mode 100644 include/erofs/cond.h

diff --git a/fsck/main.c b/fsck/main.c
index b2d8f1a..b74ae68 100644
--- a/fsck/main.c
+++ b/fsck/main.c
@@ -16,6 +16,7 @@
 #include "../lib/compressor.h"
 #include "../lib/liberofs_compress.h"
 #include "../lib/liberofs_sha256.h"
+#include "erofs/internal.h"
 
 static int erofsfsck_check_inode(erofs_nid_t pnid, erofs_nid_t nid);
 
@@ -526,6 +527,13 @@ static int erofs_verify_inode_data(struct erofs_inode *inode, int outfd,
 	u64 pchunk_len = 0;
 	u64 raw_size = 0, buffer_size = 0;
 	char *raw = NULL, *buffer = NULL;
+	struct z_erofs_mt_read_ctx *ctx = NULL;
+
+	if (!digest) {
+		ctx = z_erofs_mt_read_ctx_alloc(outfd, true);
+		if (!ctx)
+			return -ENOMEM;
+	}
 
 	erofs_dbg("verify data chunk of nid(%llu): type(%d)",
 		  inode->nid | 0ULL, inode->datalayout);
@@ -569,7 +577,7 @@ static int erofs_verify_inode_data(struct erofs_inode *inode, int outfd,
 						(const u8 *)zeros, chunk);
 					remain -= chunk;
 				}
-			} else if (outfd >= 0) {
+			} else if (!ctx && outfd >= 0) {
 				ret = lseek(outfd, map.m_llen, SEEK_CUR);
 				if (ret < 0) {
 					ret = -errno;
@@ -592,8 +600,9 @@ static int erofs_verify_inode_data(struct erofs_inode *inode, int outfd,
 			alloc_rawsize = map.m_plen;
 		}
 
-		if (alloc_rawsize > raw_size) {
-			char *newraw = realloc(raw, alloc_rawsize);
+		if (!ctx) {
+			if (alloc_rawsize > raw_size) {
+				char *newraw = realloc(raw, alloc_rawsize);
 
 			if (!newraw) {
 				ret = -ENOMEM;
@@ -602,9 +611,22 @@ static int erofs_verify_inode_data(struct erofs_inode *inode, int outfd,
 			raw = newraw;
 			raw_size = alloc_rawsize;
 		}
+	}
 
 		if (compressed) {
-			if (map.m_llen > buffer_size) {
+			char *c_raw = raw;
+			char *c_buf = buffer;
+
+			if (ctx) {
+				c_raw = malloc(alloc_rawsize);
+				c_buf = calloc(1, map.m_llen);
+				if (!c_raw || !c_buf) {
+					free(c_raw);
+					free(c_buf);
+					ret = -ENOMEM;
+					goto out;
+				}
+			} else if (map.m_llen > buffer_size) {
 				char *newbuffer;
 
 				buffer_size = map.m_llen;
@@ -614,45 +636,98 @@ static int erofs_verify_inode_data(struct erofs_inode *inode, int outfd,
 					goto out;
 				}
 				buffer = newbuffer;
+				c_buf = buffer;
 			}
-			ret = z_erofs_read_one_data(inode, &map, raw, buffer,
-						    0, map.m_llen, false);
-			if (ret)
+
+			ret = z_erofs_read_one_data(inode, &map, c_raw, c_buf,
+						    0, map.m_llen, false,
+						    map.m_la, ctx);
+			if (ret) {
+				if (ctx) {
+					free(c_raw);
+					free(c_buf);
+				}
 				goto out;
+			}
 
-			if (digest)
-				erofs_sha256_process(digest,
-					(const u8 *)buffer, map.m_llen);
-			if (outfd >= 0 && write(outfd, buffer, map.m_llen) < 0)
-				goto fail_eio;
+			if (!ctx) {
+				if (digest)
+					erofs_sha256_process(digest,
+						(const u8 *)c_buf, map.m_llen);
+				if (outfd >= 0 && write(outfd, c_buf, map.m_llen) < 0)
+					goto fail_eio;
+			}
 		} else {
+			char *c_raw = raw;
 			u64 p = 0;
+			erofs_off_t m_llen = map.m_llen;
+
+			if (ctx) {
+				c_raw = calloc(1, alloc_rawsize);
+				if (!c_raw) {
+					ret = -ENOMEM;
+					goto out;
+				}
+			}
 
 			do {
-				u64 count = min_t(u64, alloc_rawsize,
-						  map.m_llen);
+				u64 count = min_t(u64, alloc_rawsize, m_llen);
 
-				ret = erofs_read_one_data(inode, &map, raw, p, count);
-				if (ret)
+				ret = erofs_read_one_data(inode, &map, c_raw, p, count);
+				if (ret) {
+					if (ctx)
+						free(c_raw);
 					goto out;
+				}
 
-				if (digest)
-					erofs_sha256_process(digest,
-						(const u8 *)raw, count);
-				if (outfd >= 0 && write(outfd, raw, count) < 0)
-					goto fail_eio;
-				map.m_llen -= count;
+				if (!ctx) {
+					if (digest)
+						erofs_sha256_process(digest,
+							(const u8 *)c_raw, count);
+					if (outfd >= 0 && write(outfd, c_raw, count) < 0)
+						goto fail_eio;
+				} else if (outfd >= 0) {
+					if (pwrite(outfd, c_raw, count, map.m_la + p) < 0) {
+						free(c_raw);
+						goto fail_eio;
+					}
+				}
+				m_llen -= count;
 				p += count;
-			} while (map.m_llen);
+			} while (m_llen);
+
+			if (ctx)
+				free(c_raw);
 		}
 	}
 
+	if (ctx) {
+		int wait_err;
+
+		z_erofs_mt_read_enqueue(ctx);
+		wait_err = z_erofs_mt_read_ctx_wait(ctx);
+		if (wait_err < 0 && ret == 0)
+			ret = wait_err;
+	}
+
 	if (fsckcfg.print_comp_ratio) {
 		if (!erofs_is_packed_inode(inode))
 			fsckcfg.logical_blocks += BLK_ROUND_UP(inode->sbi, inode->i_size);
 		fsckcfg.physical_blocks += BLK_ROUND_UP(inode->sbi, pchunk_len);
 	}
 out:
+	if (ctx) {
+		/* use ftruncate to handle trailing unmapped blocks (holes) */
+		if (outfd >= 0 && ret == 0) {
+			if (ftruncate(outfd, inode->i_size) < 0) {
+				erofs_err("failed to truncate file to %llu: %d", 
+						(unsigned long long)inode->i_size, errno);
+				if (ret == 0) 
+					ret = -errno;
+			}
+		}
+		z_erofs_mt_read_ctx_free(ctx);
+	}
 	if (raw)
 		free(raw);
 	if (buffer)
@@ -1127,7 +1202,9 @@ int main(int argc, char *argv[])
 {
 	int err;
 
-	erofs_init_configure();
+	err = liberofs_global_init();
+	if (err)
+		return 1;
 
 	fsckcfg.physical_blocks = 0;
 	fsckcfg.logical_blocks = 0;
@@ -1283,7 +1360,7 @@ exit_dev_close:
 	erofs_dev_close(&g_sbi);
 exit:
 	erofs_blob_closeall(&g_sbi);
-	erofs_exit_configure();
+	liberofs_global_exit();
 	return err ? 1 : 0;
 }
 
diff --git a/include/erofs/cond.h b/include/erofs/cond.h
new file mode 100644
index 0000000..90ec838
--- /dev/null
+++ b/include/erofs/cond.h
@@ -0,0 +1,31 @@
+/* SPDX-License-Identifier: GPL-2.0+ OR Apache-2.0 */
+#ifndef __EROFS_COND_H
+#define __EROFS_COND_H
+
+#include "lock.h"
+
+#if defined(HAVE_PTHREAD_H) && defined(EROFS_MT_ENABLED)
+#include <pthread.h>
+
+typedef pthread_cond_t erofs_cond_t;
+
+static inline void erofs_cond_init(erofs_cond_t *cond)
+{
+	pthread_cond_init(cond, NULL);
+}
+#define erofs_cond_wait		pthread_cond_wait
+#define erofs_cond_signal	pthread_cond_signal
+#define erofs_cond_broadcast	pthread_cond_broadcast
+#define erofs_cond_destroy	pthread_cond_destroy
+
+#else
+typedef struct {} erofs_cond_t;
+
+static inline void erofs_cond_init(erofs_cond_t *cond) {}
+static inline int erofs_cond_wait(erofs_cond_t *cond, erofs_mutex_t *mutex) { return 0; }
+static inline int erofs_cond_signal(erofs_cond_t *cond) { return 0; }
+static inline int erofs_cond_broadcast(erofs_cond_t *cond) { return 0; }
+static inline int erofs_cond_destroy(erofs_cond_t *cond) { return 0; }
+#endif
+
+#endif
\ No newline at end of file
diff --git a/include/erofs/internal.h b/include/erofs/internal.h
index 2cc9cc8..6d3d407 100644
--- a/include/erofs/internal.h
+++ b/include/erofs/internal.h
@@ -25,6 +25,11 @@ typedef unsigned short umode_t;
 #ifdef HAVE_PTHREAD_H
 #include <pthread.h>
 #endif
+#ifdef HAVE_PTHREAD_H
+#include <pthread.h>
+#endif
+#include <erofs/lock.h>
+#include "erofs/cond.h"
 #include <stdlib.h>
 #include <string.h>
 #include "atomic.h"
@@ -62,6 +67,7 @@ struct erofs_buf {
 #define erofs_pos(sbi, nr)      ((erofs_off_t)(nr) << (sbi)->blkszbits)
 #define BLK_ROUND_UP(sbi, addr)	\
 	(roundup(addr, erofs_blksiz(sbi)) >> (sbi)->blkszbits)
+#define Z_EROFS_PCLUSTER_MAX_BATCH_SIZE 32
 
 struct erofs_buffer_head;
 struct erofs_bufmgr;
@@ -451,6 +457,17 @@ struct z_erofs_paramset {
 	char *extraopts;
 };
 
+struct z_erofs_decompress_task;
+struct z_erofs_mt_read_ctx;
+
+struct z_erofs_mt_read_ctx *z_erofs_mt_read_ctx_alloc(int outfd, bool free_out);
+int z_erofs_mt_read_ctx_wait(struct z_erofs_mt_read_ctx *ctx);
+void z_erofs_mt_read_ctx_free(struct z_erofs_mt_read_ctx *ctx);
+void z_erofs_mt_read_enqueue(struct z_erofs_mt_read_ctx *ctx);
+
+int z_erofs_mt_workers_init(void);
+void z_erofs_mt_workers_exit(void);
+
 int liberofs_global_init(void);
 void liberofs_global_exit(void);
 
@@ -487,7 +504,8 @@ int erofs_read_one_data(struct erofs_inode *inode, struct erofs_map_blocks *map,
 			char *buffer, u64 offset, size_t len);
 int z_erofs_read_one_data(struct erofs_inode *inode,
 			struct erofs_map_blocks *map, char *raw, char *buffer,
-			erofs_off_t skip, erofs_off_t length, bool trimmed);
+			erofs_off_t skip, erofs_off_t length, bool trimmed,
+			erofs_off_t out_offset, struct z_erofs_mt_read_ctx *ctx);
 void *erofs_read_metadata(struct erofs_sb_info *sbi, erofs_nid_t nid,
 			  erofs_off_t *offset, int *lengthp);
 int z_erofs_parse_cfgs(struct erofs_sb_info *sbi, struct erofs_super_block *dsb);
diff --git a/include/erofs/lock.h b/include/erofs/lock.h
index 884f23e..e08b76e 100644
--- a/include/erofs/lock.h
+++ b/include/erofs/lock.h
@@ -15,6 +15,7 @@ static inline void erofs_mutex_init(erofs_mutex_t *lock)
 }
 #define erofs_mutex_lock	pthread_mutex_lock
 #define erofs_mutex_unlock	pthread_mutex_unlock
+#define erofs_mutex_destroy	pthread_mutex_destroy
 
 #define EROFS_DEFINE_MUTEX(lock)	\
 	erofs_mutex_t lock = PTHREAD_MUTEX_INITIALIZER
@@ -35,6 +36,7 @@ typedef struct {} erofs_mutex_t;
 static inline void erofs_mutex_init(erofs_mutex_t *lock) {}
 static inline void erofs_mutex_lock(erofs_mutex_t *lock) {}
 static inline void erofs_mutex_unlock(erofs_mutex_t *lock) {}
+static inline void erofs_mutex_destroy(erofs_mutex_t *lock) {}
 
 #define EROFS_DEFINE_MUTEX(lock)	\
 	erofs_mutex_t lock = {}
diff --git a/lib/data.c b/lib/data.c
index 1bb9269..6bcd5b1 100644
--- a/lib/data.c
+++ b/lib/data.c
@@ -9,6 +9,228 @@
 #include "erofs/trace.h"
 #include "erofs/decompress.h"
 #include "liberofs_fragments.h"
+#include "erofs/lock.h"
+
+#ifdef EROFS_MT_ENABLED
+#include "erofs/config.h"
+
+struct z_erofs_percpu_queue {
+	erofs_mutex_t lock;
+	erofs_cond_t cond;
+	struct z_erofs_decompress_task *head;
+	struct z_erofs_decompress_task *tail;
+	int pending_tasks;
+	bool shutdown;
+	pthread_t thread;
+} __attribute__((aligned(64)));
+
+static struct z_erofs_percpu_queue *worker_queues;
+static int num_workers;
+#endif
+
+struct z_erofs_mt_read_ctx {
+	erofs_mutex_t lock;
+	erofs_cond_t cond;
+	int pending_tasks;
+	int final_err;
+	int outfd;
+	bool free_out;
+	struct z_erofs_decompress_task *current_task;
+};
+
+struct z_erofs_mt_read_ctx *z_erofs_mt_read_ctx_alloc(int outfd, bool free_out)
+{
+	struct z_erofs_mt_read_ctx *ctx = calloc(1, sizeof(*ctx));
+	if (!ctx)
+		return NULL;
+
+	erofs_mutex_init(&ctx->lock);
+	erofs_cond_init(&ctx->cond);
+	ctx->outfd = outfd;
+	ctx->free_out = free_out;
+	return ctx;
+}
+
+int z_erofs_mt_read_ctx_wait(struct z_erofs_mt_read_ctx *ctx)
+{
+	int err;
+	if (!ctx)
+		return 0;
+
+	erofs_mutex_lock(&ctx->lock);
+	while (ctx->pending_tasks > 0)
+		erofs_cond_wait(&ctx->cond, &ctx->lock);
+	err = ctx->final_err;
+	erofs_mutex_unlock(&ctx->lock);
+
+	return err;
+}
+
+void z_erofs_mt_read_ctx_free(struct z_erofs_mt_read_ctx *ctx)
+{
+	if (!ctx)
+		return;
+	
+	erofs_mutex_destroy(&ctx->lock);
+	erofs_cond_destroy(&ctx->cond);
+	free(ctx);
+}
+
+struct z_erofs_decompress_item {
+	struct z_erofs_decompress_req req;
+	char *raw_buf;
+	char *out_buf;
+	erofs_off_t out_offset;
+	unsigned int out_length;
+};
+
+struct z_erofs_decompress_task {
+	struct z_erofs_decompress_task *next;
+	struct z_erofs_mt_read_ctx *ctx;
+	struct z_erofs_decompress_item items[Z_EROFS_PCLUSTER_MAX_BATCH_SIZE];
+	unsigned int nr_reqs;
+};
+
+static void z_erofs_process_task(struct z_erofs_decompress_task *task)
+{
+	struct z_erofs_mt_read_ctx *ctx = task->ctx;
+	int i, ret = 0, first_err = 0;
+
+	for (i = 0; i < task->nr_reqs; ++i) {
+		struct z_erofs_decompress_item *item = &task->items[i];
+		ret = z_erofs_decompress(&item->req);
+
+		if (ret >= 0 && ctx && ctx->outfd >= 0) {
+			if (pwrite(ctx->outfd, item->out_buf,
+				   item->out_length, item->out_offset) < 0)
+				ret = -errno;
+		}
+
+		if (ret < 0 && !first_err)
+			first_err = ret;
+
+		free(item->raw_buf);
+		if (ctx && ctx->free_out)
+			free(item->out_buf);
+	}
+
+	if (ctx) {
+		erofs_mutex_lock(&ctx->lock);
+		if (first_err < 0 && !ctx->final_err)
+			ctx->final_err = first_err;
+		ctx->pending_tasks--;
+		if (!ctx->pending_tasks)
+			erofs_cond_signal(&ctx->cond);
+		erofs_mutex_unlock(&ctx->lock);
+	}
+	free(task);
+}
+
+#ifdef EROFS_MT_ENABLED
+static void *z_erofs_worker_thread(void *arg)
+{
+	struct z_erofs_percpu_queue *q = arg;
+
+	while (1) {
+		erofs_mutex_lock(&q->lock);
+		while (!q->head && !q->shutdown)
+			erofs_cond_wait(&q->cond, &q->lock);
+
+		if (q->shutdown && !q->head) {
+			erofs_mutex_unlock(&q->lock);
+			break;
+		}
+
+		struct z_erofs_decompress_task *tasks = q->head;
+		q->head = q->tail = NULL;
+		q->pending_tasks = 0;
+		erofs_mutex_unlock(&q->lock);
+
+		while (tasks) {
+			struct z_erofs_decompress_task *task = tasks;
+			tasks = tasks->next;
+			z_erofs_process_task(task);
+		}
+	}
+	return NULL;
+}
+
+int z_erofs_mt_workers_init(void)
+{
+	int i;
+	num_workers = erofs_get_available_processors();
+	if (num_workers < 1)
+		num_workers = 1;
+	
+	worker_queues = calloc(num_workers, sizeof(*worker_queues));
+	if (!worker_queues)
+		return -ENOMEM;
+
+	for (i = 0; i < num_workers; ++i) {
+		erofs_mutex_init(&worker_queues[i].lock);
+		erofs_cond_init(&worker_queues[i].cond);
+		worker_queues[i].head = worker_queues[i].tail = NULL;
+		worker_queues[i].pending_tasks = 0;
+		worker_queues[i].shutdown = false;
+		pthread_create(&worker_queues[i].thread, NULL, z_erofs_worker_thread, &worker_queues[i]);
+	}
+	return 0;
+}
+
+void z_erofs_mt_workers_exit(void)
+{
+	int i;
+	if (!worker_queues) return;
+	for (i = 0; i < num_workers; ++i) {
+		erofs_mutex_lock(&worker_queues[i].lock);
+		worker_queues[i].shutdown = true;
+		erofs_cond_signal(&worker_queues[i].cond);
+		erofs_mutex_unlock(&worker_queues[i].lock);
+
+		pthread_join(worker_queues[i].thread, NULL);
+		erofs_cond_destroy(&worker_queues[i].cond);
+		erofs_mutex_destroy(&worker_queues[i].lock);
+	}
+	free(worker_queues);
+	worker_queues = NULL;
+}
+#endif
+
+void z_erofs_mt_read_enqueue(struct z_erofs_mt_read_ctx *ctx)
+{
+#ifdef EROFS_MT_ENABLED
+	static int next_worker = 0;
+#endif
+
+	if (!ctx || !ctx->current_task)
+		return;
+
+	struct z_erofs_decompress_task *task = ctx->current_task;
+	ctx->current_task = NULL;
+
+#ifdef EROFS_MT_ENABLED
+	if (num_workers > 0) {
+		int target = next_worker;
+		next_worker = (next_worker + 1) % num_workers;
+		struct z_erofs_percpu_queue *q = &worker_queues[target];
+
+		erofs_mutex_lock(&q->lock);
+		task->next = NULL;
+		if (!q->tail) {
+			q->head = q->tail = task;
+		} else {
+			q->tail->next = task;
+			q->tail = task;
+		}
+		q->pending_tasks++;
+		erofs_cond_signal(&q->cond);
+		erofs_mutex_unlock(&q->lock);
+		return;
+	}
+#endif
+	task->next = NULL;
+	z_erofs_process_task(task);
+}
 
 void *erofs_bread(struct erofs_buf *buf, erofs_off_t offset, bool need_kmap)
 {
@@ -277,20 +499,29 @@ static int erofs_read_raw_data(struct erofs_inode *inode, char *buffer,
 
 int z_erofs_read_one_data(struct erofs_inode *inode,
 			struct erofs_map_blocks *map, char *raw, char *buffer,
-			erofs_off_t skip, erofs_off_t length, bool trimmed)
+			erofs_off_t skip, erofs_off_t length, bool trimmed,
+			erofs_off_t out_offset, struct z_erofs_mt_read_ctx *ctx)
 {
 	struct erofs_sb_info *sbi = inode->sbi;
 	struct erofs_map_dev mdev;
-	int ret = 0;
+	struct z_erofs_decompress_task *task;
+	struct z_erofs_decompress_item *item;
+	int ret = 0, idx, batch_limit;
 
 	if (map->m_flags & __EROFS_MAP_FRAGMENT) {
 		if (__erofs_unlikely(inode->nid == sbi->packed_nid)) {
 			erofs_err("fragment should not exist in the packed inode %llu",
-				  sbi->packed_nid | 0ULL);
-			return -EFSCORRUPTED;
+				sbi->packed_nid | 0ULL);
+			ret = -EFSCORRUPTED;
+			goto err_out;
+		}
+		ret = erofs_packedfile_read(sbi, buffer, length - skip,
+					inode->fragmentoff + skip);
+		if (ret >= 0 && ctx && ctx->outfd >= 0) {
+			if (pwrite(ctx->outfd, buffer, length - skip, out_offset) < 0)
+				ret = -errno;
 		}
-		return erofs_packedfile_read(sbi, buffer, length - skip,
-				   inode->fragmentoff + skip);
+		goto err_out;
 	}
 
 	/* no device id here, thus it will always succeed */
@@ -300,31 +531,84 @@ int z_erofs_read_one_data(struct erofs_inode *inode,
 	ret = erofs_map_dev(sbi, &mdev);
 	if (ret) {
 		DBG_BUGON(1);
-		return ret;
+		goto err_out;
 	}
 
 	ret = erofs_dev_read(sbi, mdev.m_deviceid, raw, mdev.m_pa, map->m_plen);
 	if (ret < 0)
+		goto err_out;
+
+	if (!ctx) {
+		ret = z_erofs_decompress(&(struct z_erofs_decompress_req) {
+				.sbi = sbi,
+				.in = raw,
+				.out = buffer,
+				.decodedskip = skip,
+				.interlaced_offset =
+					map->m_algorithmformat == Z_EROFS_COMPRESSION_INTERLACED ?
+						erofs_blkoff(sbi, map->m_la) : 0,
+				.inputsize = map->m_plen,
+				.decodedlength = length,
+				.alg = map->m_algorithmformat,
+				.partial_decoding = trimmed ? true :
+					!(map->m_flags & EROFS_MAP_FULL_MAPPED) ||
+						(map->m_flags & EROFS_MAP_PARTIAL_REF),
+				});
 		return ret;
+	}
 
-	ret = z_erofs_decompress(&(struct z_erofs_decompress_req) {
-			.sbi = sbi,
-			.in = raw,
-			.out = buffer,
-			.decodedskip = skip,
-			.interlaced_offset =
-				map->m_algorithmformat == Z_EROFS_COMPRESSION_INTERLACED ?
-					erofs_blkoff(sbi, map->m_la) : 0,
-			.inputsize = map->m_plen,
-			.decodedlength = length,
-			.alg = map->m_algorithmformat,
-			.partial_decoding = trimmed ? true :
-				!(map->m_flags & EROFS_MAP_FULL_MAPPED) ||
-					(map->m_flags & EROFS_MAP_PARTIAL_REF),
-			 });
-	if (ret < 0)
-		return ret;
+	task = ctx->current_task;
+	if (!task) {
+		task = calloc(1, sizeof(*task));
+		if (!task) {
+			ret = -ENOMEM;
+			goto err_out;
+		}
+		task->ctx = ctx;
+		ctx->current_task = task;
+
+		erofs_mutex_lock(&ctx->lock);
+		ctx->pending_tasks++;
+		erofs_mutex_unlock(&ctx->lock);
+	}
+
+	idx = task->nr_reqs++;
+	item = &task->items[idx];
+
+	item->req = (struct z_erofs_decompress_req) {
+		.sbi = sbi,
+		.in = raw,
+		.out = buffer,
+		.decodedskip = skip,
+		.interlaced_offset =
+			map->m_algorithmformat == Z_EROFS_COMPRESSION_INTERLACED ?
+				erofs_blkoff(sbi, map->m_la) : 0,
+		.inputsize = map->m_plen,
+		.decodedlength = length,
+		.alg = map->m_algorithmformat,
+		.partial_decoding = trimmed ? true :
+			!(map->m_flags & EROFS_MAP_FULL_MAPPED) ||
+				(map->m_flags & EROFS_MAP_PARTIAL_REF),
+	};
+	item->raw_buf = raw;
+	item->out_buf = buffer;
+	item->out_offset = out_offset;
+	item->out_length = length;
+
+	batch_limit = (map->m_algorithmformat == Z_EROFS_COMPRESSION_LZ4) ?
+					Z_EROFS_PCLUSTER_MAX_BATCH_SIZE : 8;
+
+	if (task->nr_reqs >= batch_limit)
+		z_erofs_mt_read_enqueue(ctx);
 	return 0;
+
+err_out:
+	if (ctx) {
+		if (ctx->free_out)
+			free(buffer);
+		free(raw);
+	}
+	return ret;
 }
 
 static int z_erofs_read_data(struct erofs_inode *inode, char *buffer,
@@ -387,7 +671,7 @@ static int z_erofs_read_data(struct erofs_inode *inode, char *buffer,
 		}
 
 		ret = z_erofs_read_one_data(inode, &map, raw,
-				buffer + end - offset, skip, length, trimmed);
+				buffer + end - offset, skip, length, trimmed, 0, NULL);
 		if (ret < 0)
 			break;
 	}
diff --git a/lib/global.c b/lib/global.c
index 938aa0a..8bb9da0 100644
--- a/lib/global.c
+++ b/lib/global.c
@@ -12,6 +12,7 @@
 #include "erofs/err.h"
 #include "erofs/config.h"
 #include "liberofs_compress.h"
+#include "erofs/internal.h"
 
 static EROFS_DEFINE_MUTEX(erofs_global_mutex);
 #ifdef HAVE_LIBCURL
@@ -27,6 +28,11 @@ int liberofs_global_init(void)
 #ifdef S3EROFS_ENABLED
 	xmlInitParser();
 #endif
+#ifdef EROFS_MT_ENABLED
+	err = z_erofs_mt_workers_init();
+	if (err)
+		goto out_unlock;
+#endif
 #ifdef HAVE_LIBCURL
 	if (!erofs_global_curl_initialized) {
 		if (curl_global_init(CURL_GLOBAL_DEFAULT) != CURLE_OK) {
@@ -35,8 +41,8 @@ int liberofs_global_init(void)
 		}
 		erofs_global_curl_initialized = true;
 	}
-out_unlock:
 #endif
+out_unlock:
 	erofs_mutex_unlock(&erofs_global_mutex);
 	return err;
 }
@@ -45,6 +51,9 @@ void liberofs_global_exit(void)
 {
 	erofs_mutex_lock(&erofs_global_mutex);
 	z_erofs_mt_global_exit();
+#ifdef EROFS_MT_ENABLED
+	z_erofs_mt_workers_exit();
+#endif
 #ifdef HAVE_LIBCURL
 	if (erofs_global_curl_initialized) {
 		curl_global_cleanup();
-- 
2.53.0



  parent reply	other threads:[~2026-07-13  7:03 UTC|newest]

Thread overview: 22+ messages / expand[flat|nested]  mbox.gz  Atom feed  top
2026-06-21 12:01 [PATCH 0/2] fsck.erofs: implement multi-threaded extraction Nithurshen
2026-06-21 12:01 ` [PATCH 1/2] fsck.erofs: add multi-threaded decompression Nithurshen
2026-06-22  2:08   ` Gao Xiang
2026-06-29  9:55   ` [PATCH 1/2 v2] " Nithurshen
2026-07-05 14:00     ` Gao Xiang
2026-07-06  6:05   ` [PATCH 1/2 v3] " Nithurshen
2026-07-06  6:10     ` Nithurshen
2026-07-08  2:31       ` Gao Xiang
2026-07-13  7:02     ` Nithurshen [this message]
2026-06-21 12:01 ` [PATCH 2/2] fsck.erofs: implement concurrent directory traversal Nithurshen
2026-07-05 14:05   ` Gao Xiang
2026-07-14  1:10   ` [PATCH v5] fsck.erofs: add multi-threaded decompression Nithurshen
2026-07-14  1:18     ` Nithurshen Karthikeyan
2026-07-26 11:01     ` [PATCH v6] " Nithurshen
2026-07-26 11:02       ` Nithurshen
2026-07-24 14:35   ` [PATCH v2 2/2] fsck.erofs: implement concurrent directory traversal Nithurshen
2026-07-24 14:47     ` Nithurshen
2026-07-24 15:28       ` Gao Xiang
2026-07-24 15:32         ` Nithurshen Karthikeyan
2026-07-24 15:45           ` Gao Xiang
2026-07-24 15:53     ` Gao Xiang
2026-07-25  1:26   ` [PATCH v3 " Nithurshen

Reply instructions:

You may reply publicly to this message via plain-text email
using any one of the following methods:

* Save the following mbox file, import it into your mail client,
  and reply-to-all from there: mbox

  Avoid top-posting and favor interleaved quoting:
  https://en.wikipedia.org/wiki/Posting_style#Interleaved_style

* Reply using the --to, --cc, and --in-reply-to
  switches of git-send-email(1):

  git send-email \
    --in-reply-to=20260713070249.75652-1-nithurshen.dev@gmail.com \
    --to=nithurshen.dev@gmail.com \
    --cc=hsiangkao@linux.alibaba.com \
    --cc=linux-erofs@lists.ozlabs.org \
    --cc=xiang@kernel.org \
    /path/to/YOUR_REPLY

  https://kernel.org/pub/software/scm/git/docs/git-send-email.html

* If your mail client supports setting the In-Reply-To header
  via mailto: links, try the mailto: link
Be sure your reply has a Subject: header at the top and a blank line before the message body.
This is a public inbox, see mirroring instructions
for how to clone and mirror all data and code used for this inbox