git/list[1] front-page[2] threads[3] people[4] search[5] about
 

[PATCH v3 05/12] object-store: allow threaded access to object reading

From
Matheus Tavares <matheus.bernardino@usp.br>
Date
Jan 16, 2020, 02:39 UTC
Message-ID
<b72e90f229dbf7d5be016fd6251a9b3ef76f2431.1579141989.git.matheus.bernardino@usp.br>
In-Reply-To
<cover.1579141989.git.matheus.bernardino@usp.br>

Allow object reading to be performed by multiple threads protecting it with an internal lock, the obj_read_mutex. The lock usage can be toggled with enable_obj_read_lock() and disable_obj_read_lock(). Currently, the functions which can be safely called in parallel are: read_object_file_extended(), repo_read_object_file(), read_object_file(), read_object_with_reference(), read_object(), oid_object_info() and oid_object_info_extended(). It's also possible to use obj_read_lock() and obj_read_unlock() to protect other sections that cannot execute in parallel with object reading.

Probably there are many spots in the functions listed above that could be executed unlocked (and thus, in parallel). But, for now, we are most interested in allowing parallel access to zlib inflation. This is one of the sections where object reading spends most of the time in (e.g. up to one-third of git-grep's execution time in the chromium repo corresponds to inflation) and it's already thread-safe. So, to take advantage of that, the obj_read_mutex is released when calling git_inflate() and re-acquired right after, for every calling spot in oid_object_info_extended()'s call chain. We may refine this lock to also exploit other possible parallel spots in the future, but for now, threaded zlib inflation should already give great speedups for threaded object reading callers.

Note that add_delta_base_cache() was also modified to skip adding already present entries to the cache. This wasn't possible before, but it would be now, with the parallel inflation. Take for example the following situation, where two threads - A and B - are executing the code at unpack_entry():

1. Thread A is performing the decompression of a base O (which is not
   yet in the cache) at PHASE II. Thread B is simultaneously trying to
   unpack O, but just starting at PHASE I.
2. Since O is not yet in the cache, B will go to PHASE II to also
   perform the decompression.
3. When they finish decompressing, one of them will get the object
   reading mutex and go to PHASE III while the other waits for the
   mutex. Let’s say A got the mutex first.
4. Thread A will add O to the cache, go throughout the rest of PHASE III
   and return.
5. Thread B gets the mutex, also add O to the cache (if the check wasn't
   there) and returns.

Finally, it is also important to highlight that the object reading lock can only ensure thread-safety in the mentioned functions thanks to two complementary mechanisms: the use of 'struct raw_object_store's replace_mutex, which guards sections in the object reading machinery that would otherwise be thread-unsafe; and the 'struct pack_window's inuse_cnt, which protects window reading operations (such as the one performed during the inflation of a packed object), allowing them to execute without the acquisition of the obj_read_mutex.

Signed-off-by: Matheus Tavares <matheus.bernardino@usp.br>
---
 object-store.h | 35 +++++++++++++++++++++++++++++++
 packfile.c     | 32 ++++++++++++++++++++++++++++
 sha1-file.c    | 57 +++++++++++++++++++++++++++++++++++++++++++++-----
 3 files changed, 119 insertions(+), 5 deletions(-)
diff --git a/object-store.h b/object-store.h
index 33739c9dee..7c80e0d64c 100644
--- a/object-store.h
+++ b/object-store.h
@@ -6,6 +6,7 @@
 #include "list.h"
 #include "sha1-array.h"
 #include "strbuf.h"
+#include "thread-utils.h"
 
 struct object_directory {
 	struct object_directory *next;
@@ -251,6 +252,40 @@ int has_loose_object_nonlocal(const struct object_id *);
 
 void assert_oid_type(const struct object_id *oid, enum object_type expect);
 
+/*
+ * Enabling the object read lock allows multiple threads to safely call the
+ * following functions in parallel: repo_read_object_file(), read_object_file(),
+ * read_object_file_extended(), read_object_with_reference(), read_object(),
+ * oid_object_info() and oid_object_info_extended().
+ *
+ * obj_read_lock() and obj_read_unlock() may also be used to protect other
+ * section which cannot execute in parallel with object reading. Since the used
+ * lock is a recursive mutex, these sections can even contain calls to object
+ * reading functions. However, beware that in these cases zlib inflation won't
+ * be performed in parallel, losing performance.
+ *
+ * TODO: oid_object_info_extended()'s call stack has a recursive behavior. If
+ * any of its callees end up calling it, this recursive call won't benefit from
+ * parallel inflation.
+ */
+void enable_obj_read_lock(void);
+void disable_obj_read_lock(void);
+
+extern int obj_read_use_lock;
+extern pthread_mutex_t obj_read_mutex;
+
+static inline void obj_read_lock(void)
+{
+	if(obj_read_use_lock)
+		pthread_mutex_lock(&obj_read_mutex);
+}
+
+static inline void obj_read_unlock(void)
+{
+	if(obj_read_use_lock)
+		pthread_mutex_unlock(&obj_read_mutex);
+}
+
 struct object_info {
 	/* Request */
 	enum object_type *typep;
diff --git a/packfile.c b/packfile.c
index 7e7c04e4d8..24a73fc33a 100644
--- a/packfile.c
+++ b/packfile.c
@@ -1086,7 +1086,23 @@ unsigned long get_size_from_delta(struct packed_git *p,
 	do {
 		in = use_pack(p, w_curs, curpos, &stream.avail_in);
 		stream.next_in = in;
+		/*
+		 * Note: the window section returned by use_pack() must be
+		 * available throughout git_inflate()'s unlocked execution. To
+		 * ensure no other thread will modify the window in the
+		 * meantime, we rely on the packed_window.inuse_cnt. This
+		 * counter is incremented before window reading and checked
+		 * before window disposal.
+		 *
+		 * Other worrying sections could be the call to close_pack_fd(),
+		 * which can close packs even with in-use windows, and to
+		 * reprepare_packed_git(). Regarding the former, mmap doc says:
+		 * "closing the file descriptor does not unmap the region". And
+		 * for the latter, it won't re-open already available packs.
+		 */
+		obj_read_unlock();
 		st = git_inflate(&stream, Z_FINISH);
+		obj_read_lock();
 		curpos += stream.next_in - in;
 	} while ((st == Z_OK || st == Z_BUF_ERROR) &&
 		 stream.total_out < sizeof(delta_head));
@@ -1445,6 +1461,14 @@ static void add_delta_base_cache(struct packed_git *p, off_t base_offset,
 	struct delta_base_cache_entry *ent = xmalloc(sizeof(*ent));
 	struct list_head *lru, *tmp;
 
+	/*
+	 * Check required to avoid redundant entries when more than one thread
+	 * is unpacking the same object, in unpack_entry() (since its phases I
+	 * and III might run concurrently across multiple threads).
+	 */
+	if (in_delta_base_cache(p, base_offset))
+		return;
+
 	delta_base_cached += base_size;
 
 	list_for_each_safe(lru, tmp, &delta_base_cache_lru) {
@@ -1574,7 +1598,15 @@ static void *unpack_compressed_entry(struct packed_git *p,
 	do {
 		in = use_pack(p, w_curs, curpos, &stream.avail_in);
 		stream.next_in = in;
+		/*
+		 * Note: we must ensure the window section returned by
+		 * use_pack() will be available throughout git_inflate()'s
+		 * unlocked execution. Please refer to the comment at
+		 * get_size_from_delta() to see how this is done.
+		 */
+		obj_read_unlock();
 		st = git_inflate(&stream, Z_FINISH);
+		obj_read_lock();
 		if (!stream.avail_out)
 			break; /* the payload is larger than it should be */
 		curpos += stream.next_in - in;
diff --git a/sha1-file.c b/sha1-file.c
index 188de57634..9dc0649748 100644
--- a/sha1-file.c
+++ b/sha1-file.c
@@ -1147,6 +1147,8 @@ static int unpack_loose_short_header(git_zstream *stream,
 				     unsigned char *map, unsigned long mapsize,
 				     void *buffer, unsigned long bufsiz)
 {
+	int ret;
+
 	/* Get the data stream */
 	memset(stream, 0, sizeof(*stream));
 	stream->next_in = map;
@@ -1155,7 +1157,11 @@ static int unpack_loose_short_header(git_zstream *stream,
 	stream->avail_out = bufsiz;
 
 	git_inflate_init(stream);
-	return git_inflate(stream, 0);
+	obj_read_unlock();
+	ret = git_inflate(stream, 0);
+	obj_read_lock();
+
+	return ret;
 }
 
 int unpack_loose_header(git_zstream *stream,
@@ -1200,7 +1206,9 @@ static int unpack_loose_header_to_strbuf(git_zstream *stream, unsigned char *map
 	stream->avail_out = bufsiz;
 
 	do {
+		obj_read_unlock();
 		status = git_inflate(stream, 0);
+		obj_read_lock();
 		strbuf_add(header, buffer, stream->next_out - (unsigned char *)buffer);
 		if (memchr(buffer, '\0', stream->next_out - (unsigned char *)buffer))
 			return 0;
@@ -1240,8 +1248,11 @@ static void *unpack_loose_rest(git_zstream *stream,
 		 */
 		stream->next_out = buf + bytes;
 		stream->avail_out = size - bytes;
-		while (status == Z_OK)
+		while (status == Z_OK) {
+			obj_read_unlock();
 			status = git_inflate(stream, Z_FINISH);
+			obj_read_lock();
+		}
 	}
 	if (status == Z_STREAM_END && !stream->avail_in) {
 		git_inflate_end(stream);
@@ -1411,10 +1422,32 @@ static int loose_object_info(struct repository *r,
 	return (status < 0) ? status : 0;
 }
 
+int obj_read_use_lock = 0;
+pthread_mutex_t obj_read_mutex;
+
+void enable_obj_read_lock(void)
+{
+	if (obj_read_use_lock)
+		return;
+
+	obj_read_use_lock = 1;
+	init_recursive_mutex(&obj_read_mutex);
+}
+
+void disable_obj_read_lock(void)
+{
+	if (!obj_read_use_lock)
+		return;
+
+	obj_read_use_lock = 0;
+	pthread_mutex_destroy(&obj_read_mutex);
+}
+
 int fetch_if_missing = 1;
 
-int oid_object_info_extended(struct repository *r, const struct object_id *oid,
-			     struct object_info *oi, unsigned flags)
+static int do_oid_object_info_extended(struct repository *r,
+				       const struct object_id *oid,
+				       struct object_info *oi, unsigned flags)
 {
 	static struct object_info blank_oi = OBJECT_INFO_INIT;
 	struct pack_entry e;
@@ -1422,6 +1455,7 @@ int oid_object_info_extended(struct repository *r, const struct object_id *oid,
 	const struct object_id *real = oid;
 	int already_retried = 0;
 
+
 	if (flags & OBJECT_INFO_LOOKUP_REPLACE)
 		real = lookup_replace_object(r, oid);
 
@@ -1497,7 +1531,7 @@ int oid_object_info_extended(struct repository *r, const struct object_id *oid,
 	rtype = packed_object_info(r, e.p, e.offset, oi);
 	if (rtype < 0) {
 		mark_bad_packed_object(e.p, real->hash);
-		return oid_object_info_extended(r, real, oi, 0);
+		return do_oid_object_info_extended(r, real, oi, 0);
 	} else if (oi->whence == OI_PACKED) {
 		oi->u.packed.offset = e.offset;
 		oi->u.packed.pack = e.p;
@@ -1508,6 +1542,17 @@ int oid_object_info_extended(struct repository *r, const struct object_id *oid,
 	return 0;
 }
 
+int oid_object_info_extended(struct repository *r, const struct object_id *oid,
+			     struct object_info *oi, unsigned flags)
+{
+	int ret;
+	obj_read_lock();
+	ret = do_oid_object_info_extended(r, oid, oi, flags);
+	obj_read_unlock();
+	return ret;
+}
+
+
 /* returns enum object_type or negative */
 int oid_object_info(struct repository *r,
 		    const struct object_id *oid,
@@ -1580,6 +1625,7 @@ void *read_object_file_extended(struct repository *r,
 	if (data)
 		return data;
 
+	obj_read_lock();
 	if (errno && errno != ENOENT)
 		die_errno(_("failed to read object %s"), oid_to_hex(oid));
 
@@ -1595,6 +1641,7 @@ void *read_object_file_extended(struct repository *r,
 	if ((p = has_packed_and_bad(r, repl->hash)) != NULL)
 		die(_("packed object %s (stored in %s) is corrupt"),
 		    oid_to_hex(repl), p->pack_name);
+	obj_read_unlock();
 
 	return NULL;
 }
-- 
2.24.1
Previous: Matheus TavaresNext: Matheus Tavares
Message 33 of 47 in “grep: re-enable threads when cached, w/ parallel inflation”
  1. Matheus TavaresAug 10, 2019
  2. [GSoC][PATCH 1/4] object-store: add lock to read_object_file_extended()Matheus Tavares, Aug 10, 2019
  3. [GSoC][PATCH 2/4] grep: allow locks to be enabled individuallyMatheus Tavares, Aug 10, 2019
  4. [GSoC][PATCH 3/4] grep: disable grep_read_mutex when possibleMatheus Tavares, Aug 10, 2019
  5. [GSoC][PATCH 4/4] grep: re-enable threads in some non-worktree casesMatheus Tavares, Aug 10, 2019
  6. 00/11 grep: improve threading and fix race conditionsMatheus Tavares, Sep 30, 2019
  7. 01/11 grep: fix race conditions on userdiff callsMatheus Tavares, Sep 30, 2019
  8. 02/11 grep: fix race conditions at grep_submodule()Matheus Tavares, Sep 30, 2019
  9. 03/11 grep: fix racy calls in grep_objects()Matheus Tavares, Sep 30, 2019
  10. 04/11 replace-object: make replace operations thread-safeMatheus Tavares, Sep 30, 2019
  11. 05/11 object-store: allow threaded access to object readingMatheus Tavares, Sep 30, 2019
  12. Jonathan TanNov 12, 2019
  13. Jeff KingNov 13, 2019
  14. Matheus Tavares BernardinoNov 14, 2019
  15. Jeff KingNov 14, 2019
  16. Jonathan TanNov 14, 2019
  17. Jeff KingNov 15, 2019
  18. Matheus Tavares BernardinoDec 19, 2019
  19. Matheus Tavares BernardinoJan 9, 2020
  20. Christian CouderJan 10, 2020
  21. 06/11 grep: replace grep_read_mutex by internal obj read lockMatheus Tavares, Sep 30, 2019
  22. squash! grep: replace grep_read_mutex by internal obj read lockMatheus Tavares, Oct 1, 2019
  23. 07/11 submodule-config: add skip_if_read option to repo_read_gitmodules()Matheus Tavares, Sep 30, 2019
  24. 08/11 grep: allow submodule functions to run in parallelMatheus Tavares, Sep 30, 2019
  25. 09/11 grep: protect packed_git [re-]initializationMatheus Tavares, Sep 30, 2019
  26. 10/11 grep: re-enable threads in non-worktree caseMatheus Tavares, Sep 30, 2019
  27. 11/11 grep: move driver pre-load out of critical sectionMatheus Tavares, Sep 30, 2019
  28. 00/12 grep: improve threading and fix race conditionsMatheus Tavares, Jan 16, 2020
  29. 01/12 grep: fix race conditions on userdiff callsMatheus Tavares, Jan 16, 2020
  30. 02/12 grep: fix race conditions at grep_submodule()Matheus Tavares, Jan 16, 2020
  31. 03/12 grep: fix racy calls in grep_objects()Matheus Tavares, Jan 16, 2020
  32. 04/12 replace-object: make replace operations thread-safeMatheus Tavares, Jan 16, 2020
  33. 05/12 object-store: allow threaded access to object readingMatheus Tavares, Jan 16, 2020
  34. 06/12 grep: replace grep_read_mutex by internal obj read lockMatheus Tavares, Jan 16, 2020
  35. 07/12 submodule-config: add skip_if_read option to repo_read_gitmodules()Matheus Tavares, Jan 16, 2020
  36. 08/12 grep: allow submodule functions to run in parallelMatheus Tavares, Jan 16, 2020
  37. SZEDER GáborJan 29, 2020
  38. Junio C HamanoJan 29, 2020
  39. Junio C HamanoJan 29, 2020
  40. Matheus Tavares BernardinoJan 29, 2020
  41. Philippe BlainJan 30, 2020
  42. 09/12 grep: protect packed_git [re-]initializationMatheus Tavares, Jan 16, 2020
  43. 10/12 grep: re-enable threads in non-worktree caseMatheus Tavares, Jan 16, 2020
  44. 11/12 grep: move driver pre-load out of critical sectionMatheus Tavares, Jan 16, 2020
  45. 12/12 grep: use no. of cores as the default no. of threadsMatheus Tavares, Jan 16, 2020
  46. Victor LeschukJan 16, 2020
  47. Matheus TavaresJan 16, 2020

Read the whole thread, see it on lore, or plain text.

$ cat FOOTERMessages come from the public archive at lore.kernel.org/git, fetched every hour. The front page is chosen and written each morning by an AI editor and can be wrong; the threads themselves are the record. About and API. For agents: an MCP server at https://gitlist.dev/mcp, and any thread, story or person page as Markdown by adding .md to its URL (or sending Accept: text/markdown). Details in /llms.txt.