From 86c9829a239570d79ffffe53b189404a2c2de1e4 Mon Sep 17 00:00:00 2001 From: Caleb Gardner Date: Sat, 20 Jun 2026 06:34:33 -0500 Subject: [PATCH] Re-working multi-threaded extraction --- src/data/reader.zig | 20 +- src/extract-multi.zig | 457 +++++++++++------------------------------ src/extract-single.zig | 1 + 3 files changed, 139 insertions(+), 339 deletions(-) diff --git a/src/data/reader.zig b/src/data/reader.zig index d3b4d55..8608adf 100644 --- a/src/data/reader.zig +++ b/src/data/reader.zig @@ -27,7 +27,7 @@ frag_offset: u32 = 0, io: ?Io = null, cache: ?*Cache = null, -block: [1024 * 1024]u8 = undefined, +block_alloc: bool = false, interface: Io.Reader = .{ .buffer = &[0]u8{}, @@ -55,6 +55,10 @@ pub fn init(alloc: std.mem.Allocator, data: []u8, decomp: Decomp.Fn, block_size: .offset = data_start, }; } +pub fn deinit(self: *Reader) void { + if (self.block_alloc) + self.alloc.free(self.interface.buffer); +} pub fn addFrag(self: *Reader, frag_data: []u8, frag_offset: u32) void { self.frag_data = frag_data; self.frag_offset = frag_offset; @@ -65,6 +69,8 @@ pub fn addCache(self: *Reader, io: Io, cache: *Cache) void { } fn advance(self: *Reader) Io.Reader.Error!void { + if (self.block_alloc) self.alloc.free(self.interface.buffer); + if (self.block_idx > self.blocks.len) return error.EndOfStream; defer self.block_idx += 1; @@ -92,6 +98,8 @@ fn advance(self: *Reader) Io.Reader.Error!void { const block = self.blocks[self.block_idx]; defer self.offset += block.size; + std.debug.print("offset: {} block: {any}\n", .{ self.offset, block }); + if (block.size == 0) { self.sparse_block = true; self.interface.end = size; @@ -107,8 +115,14 @@ fn advance(self: *Reader) Io.Reader.Error!void { } if (self.cache == null) { - _ = self.decomp(self.alloc, self.data[self.offset..][0..block.size], self.block[0..size]) catch return error.ReadFailed; - self.interface.buffer = self.block[0..size]; + self.block_alloc = true; + self.interface.buffer = self.alloc.alloc(u8, size) catch return error.ReadFailed; + errdefer { + self.alloc.free(self.interface.buffer); + self.interface.buffer = &[0]u8{}; + } + + _ = self.decomp(self.alloc, self.data[self.offset..][0..block.size], self.interface.buffer) catch return error.ReadFailed; self.interface.end = size; } else { self.interface.buffer = self.cache.?.get(self.io.?, self.offset, block.size) catch return error.ReadFailed; diff --git a/src/extract-multi.zig b/src/extract-multi.zig index a1eeafb..5a4a874 100644 --- a/src/extract-multi.zig +++ b/src/extract-multi.zig @@ -12,6 +12,7 @@ const MetadataReader = @import("meta_rdr.zig"); const Superblock = @import("archive.zig").Superblock; const Cache = @import("util/cache.zig"); const XattrTable = @import("xattr.zig"); +const Atomic = std.atomic.Value; pub fn extract( alloc: std.mem.Allocator, @@ -31,138 +32,22 @@ pub fn extract( var xattr_table: XattrTable = .init(alloc, data, decomp, super.xattr_start); defer xattr_table.deinit(); - var buf: [150]ReturnUnion = undefined; - var sel: Io.Select(ReturnUnion) = .init(io, &buf); - defer while (sel.cancel()) |res| - switch (res) { - .path => |p| { - const path_return = p catch continue; - if (path_return.path.len != path.len) - alloc.free(path_return.path); - }, - else => {}, - }; - - var loop = io.async(finishLoop, .{ alloc, io, &sel, &id_table, &xattr_table, inode.hdr.num, options }); - var cache: Cache = .init(alloc, data, decomp); defer cache.deinit(); var frag_table: Lookup.Table(Lookup.FragEntry) = .init(alloc, data, decomp, super.frag_start, super.frag_count); defer frag_table.deinit(); - switch (inode.hdr.type) { - .file, .ext_file => sel.async( - .path, - extractFile, - .{ alloc, io, data, decomp, super.block_size, &cache, &frag_table, inode, path, true }, - ), - .dir, .ext_dir => sel.async( - .path, - extractDir, - .{ alloc, io, super, data, decomp, &sel, &cache, &frag_table, inode, path, true }, - ), - .symlink, .ext_symlink => sel.async( - .void, - extractSymlink, - .{ alloc, io, inode, path, true }, - ), - else => sel.async( - .path, - extractNod, - .{ alloc, inode, path, true }, - ), - } + var atomic_err: Atomic(?Error) = .init(null); - try loop.await(io); -} + var future = switch (inode.hdr.type) {}; -fn dirOrder(_: void, a: PathReturn, b: PathReturn) std.math.Order { - return std.math.order(std.mem.count(u8, a.path, "/"), std.mem.count(u8, b.path, "/")); -} -fn finishLoop(alloc: std.mem.Allocator, io: Io, sel: *Io.Select(ReturnUnion), id_table: *Lookup.Table(u16), xattr_table: *XattrTable, start_num: u32, options: ExtractionOptions) !void { - var dirs: std.PriorityDequeue(PathReturn, void, dirOrder) = .empty; - defer dirs.deinit(alloc); - errdefer while (dirs.popMax()) |d| - if (d.hdr.num != start_num) alloc.free(d.path); + _ = io; + _ = inode; + _ = options; + _ = path; - while (true) { - const value: ReturnUnion = try sel.await(); - - const path_ret = switch (value) { - .void => { - _ = sel.group.token.load(.unordered) orelse break; - continue; - }, - .path => |p| try p, - }; - - if (options.ignore_permissions and (options.ignore_xattr or path_ret.xattr_idx == null)) { - if (path_ret.hdr.num != start_num) - alloc.free(path_ret.path); - continue; - } - - if (path_ret.hdr.type == .dir or path_ret.hdr.type == .ext_dir) { - dirs.push(alloc, path_ret) catch |err| { - if (path_ret.hdr.num != start_num) - alloc.free(path_ret.path); - return err; - }; - continue; - } - defer if (path_ret.hdr.num != start_num) - alloc.free(path_ret.path); - - var file = try Io.Dir.cwd().openFile(io, path_ret.path, .{}); - defer file.close(io); - - if (!options.ignore_xattr and path_ret.xattr_idx != null) { - const xattr = try xattr_table.get(alloc, io, path_ret.xattr_idx.?); - defer xattr.deinit(alloc); - - for (xattr.kvs) |kv| { - const res = std.os.linux.fsetxattr(file.handle, kv.key, kv.value.ptr, kv.value.len, 0); - if (res != 0) - return error.SetXattrError; - } - } - if (!options.ignore_permissions) { - try file.setTimestamps(io, .{ - .modify_timestamp = .init(Io.Timestamp.fromNanoseconds(@as(i96, @intCast(path_ret.hdr.mod_time)) * std.time.ns_per_s)), - }); - try file.setPermissions(io, @enumFromInt(path_ret.hdr.permissions)); - try file.setOwner(io, try id_table.get(io, path_ret.hdr.uid_idx), try id_table.get(io, path_ret.hdr.gid_idx)); - } - - _ = sel.group.token.load(.unordered) orelse break; - } - - while (dirs.popMax()) |path_ret| { - defer if (path_ret.hdr.num != start_num) - alloc.free(path_ret.path); - - var file = try Io.Dir.cwd().openFile(io, path_ret.path, .{}); - defer file.close(io); - - if (!options.ignore_xattr and path_ret.xattr_idx != null) { - const xattr = try xattr_table.get(alloc, io, path_ret.xattr_idx.?); - defer xattr.deinit(alloc); - - for (xattr.kvs) |kv| { - const res = std.os.linux.fsetxattr(file.handle, kv.key, kv.value.ptr, kv.value.len, 0); - if (res != 0) - return error.SetXattrError; - } - } - if (!options.ignore_permissions) { - try file.setTimestamps(io, .{ - .modify_timestamp = .init(Io.Timestamp.fromNanoseconds(@as(i96, @intCast(path_ret.hdr.mod_time)) * std.time.ns_per_s)), - }); - try file.setPermissions(io, @enumFromInt(path_ret.hdr.permissions)); - try file.setOwner(io, try id_table.get(io, path_ret.hdr.uid_idx), try id_table.get(io, path_ret.hdr.gid_idx)); - } - } + return error.TODO; } fn extractDir( @@ -171,243 +56,143 @@ fn extractDir( super: Superblock, data: []u8, decomp: Decomp.Fn, - sel: *Io.Select(ReturnUnion), cache: *Cache, frag_table: *Lookup.Table(Lookup.FragEntry), + id_table: *Lookup.Table(u16), + xattr_table: *XattrTable, inode: Inode, path: []const u8, + options: ExtractionOptions, origin: bool, -) Error!PathReturn { - defer if (!origin) inode.deinit(alloc); - errdefer if (!origin) alloc.free(path); + atomic_err: *Atomic(?Error), +) error{Canceled}!void { + defer if (!origin) alloc.free(path); - var ret: PathReturn = .{ - .hdr = inode.hdr, - .path = path, + var xattr_idx: u32 = 0xFFFFFFFF; + + Io.Dir.cwd().createDirPath(io, path, .{}) catch |err| switch (err) { + error.Canceled => { + io.recancel(); + return error.Canceled; + }, + else => |e| { + atomic_err.store(e, .unordered); + return; + }, }; - try Io.Dir.cwd().createDirPath(io, path); + var group: Io.Group = blk: { + var dir: Directory = switch (inode.data) { + .dir => |d| d_blk: { + var meta: MetadataReader = .init(alloc, data, decomp, super.dir_start + d.block_start); + meta.interface.discardAll(d.block_offset) catch |err| break :blk err; - var dir: Directory = switch (inode.data) { - .dir => |d| blk: { - var meta: MetadataReader = .init(alloc, data, decomp, d.block_start + super.dir_start); - try meta.interface.discardAll(d.block_offset); + break :d_blk .init(alloc, &meta.interface, d.size) catch |err| break :blk err; + }, + .ext_dir => |d| d_blk: { + xattr_idx = d.xattr_idx; - break :blk try Directory.init(alloc, &meta.interface, d.size); - }, - .ext_dir => |d| blk: { - if (d.xattr_idx != 0xFFFFFFFF) ret.xattr_idx = d.xattr_idx; + var meta: MetadataReader = .init(alloc, data, decomp, super.dir_start + d.block_start); + meta.interface.discardAll(d.block_offset) catch |err| break :blk err; - var meta: MetadataReader = .init(alloc, data, decomp, d.block_start + super.dir_start); - try meta.interface.discardAll(d.block_offset); - - break :blk try Directory.init(alloc, &meta.interface, d.size); - }, - else => unreachable, - }; - defer dir.deinit(alloc); - - for (dir.entries) |entry| { - var new_inode: Inode = try .initEntry(alloc, data, decomp, super.inode_start, super.block_size, entry); - - const new_path = std.mem.concat(alloc, u8, &.{ path, "/", entry.name }) catch |err| { - new_inode.deinit(alloc); - return err; + break :d_blk .init(alloc, &meta.interface, d.size) catch |err| break :blk err; + }, + else => unreachable, }; + defer dir.deinit(alloc); - switch (entry.type) { - .dir => sel.async(.path, extractDir, .{ alloc, io, super, data, decomp, sel, cache, frag_table, new_inode, new_path, false }), - .file => sel.async(.path, extractFile, .{ alloc, io, data, decomp, super.block_size, cache, frag_table, new_inode, new_path, false }), - .symlink => sel.async(.void, extractSymlink, .{ alloc, io, new_inode, new_path, false }), - else => sel.async(.path, extractNod, .{ alloc, new_inode, new_path, false }), + var group: Io.Group = .init; + + for (dir.entries) |entry| { + var new_inode: Inode = .initEntry(alloc, data, decomp, super.inode_start, super.block_size, entry) catch |err| { + group.cancel(io); + break :blk err; + }; + + const new_path = std.mem.concat(alloc, u8, &.{ path, "/", entry.name }) catch |err| { + new_inode.deinit(alloc); + group.cancel(io); + break :blk err; + }; + + switch (entry.type) { + .dir => group.async(io, extractDir, .{ + alloc, + io, + super, + data, + decomp, + cache, + frag_table, + id_table, + xattr_table, + new_inode, + new_path, + options, + false, + atomic_err, + }), + .file => {}, + .symlink => {}, + else => {}, + } } - } + } catch |err| switch (err) { + error.Canceled => { + io.recancel(); + return error.Canceled; + }, + else => |e| { + atomic_err.store(e, .unordered); + return; + }, + }; - return ret; + try group.await(io); + + setMetadata(alloc, io, id_table, xattr_table, inode, path, options, xattr_idx) catch |err| switch (err) { + error.Canceled => { + io.recancel(); + return error.Canceled; + }, + else => |e| { + atomic_err.store(e, .unordered); + return; + }, + }; } -fn extractFile( +fn setMetadata( alloc: std.mem.Allocator, io: Io, - data: []u8, - decomp: Decomp.Fn, - block_size: u32, - cache: *Cache, - frag_table: *Lookup.Table(Lookup.FragEntry), + id_table: *Lookup.Table(u16), + xattr_table: *XattrTable, inode: Inode, path: []const u8, - origin: bool, -) Error!PathReturn { - defer if (!origin) inode.deinit(alloc); - errdefer if (!origin) alloc.free(path); + options: ExtractionOptions, + xattr_idx: u32, +) !void { + if (options.ignore_permissions and (options.ignore_xattr or xattr_idx == null)) return; - try io.checkCancel(); + var fil: Io.File = try Io.Dir.cwd().openFile(io, path, .{}); + defer fil.close(io); - var ret: PathReturn = .{ - .hdr = inode.hdr, - .path = path, - }; + if (!options.ignore_xattr and xattr_idx != 0xFFFFFFFF) { + const xattr = try xattr_table.get(alloc, io, xattr_idx.?); + defer xattr.deinit(alloc); - // var ext: DataExtractor = switch (inode.data) { - // .file => |f| blk: { - // var rdr: DataExtractor = .init(data, decomp, block_size, f.blocks, f.size, f.block_start); - // rdr.addCache(cache); - - // if (f.frag_idx == 0xFFFFFFFF) break :blk rdr; - - // const entry: Lookup.FragEntry = try frag_table.get(io, f.frag_idx); - // if (entry.size.uncompressed) { - // rdr.addFrag(data[entry.block_start..][0..entry.size.size], f.frag_offset); - // } else { - // rdr.addFrag(try cache.get(io, entry.block_start, entry.size.size), f.frag_offset); - // } - - // break :blk rdr; - // }, - // .ext_file => |f| blk: { - // if (f.xattr_idx != 0xFFFFFFFF) ret.xattr_idx = f.xattr_idx; - - // var rdr: DataExtractor = .init(data, decomp, block_size, f.blocks, f.size, f.block_start); - // rdr.addCache(cache); - - // if (f.frag_idx == 0xFFFFFFFF) break :blk rdr; - - // const entry: Lookup.FragEntry = try frag_table.get(io, f.frag_idx); - // if (entry.size.uncompressed) { - // rdr.addFrag(data[entry.block_start..][0..entry.size.size], f.frag_offset); - // } else { - // rdr.addFrag(try cache.get(io, entry.block_start, entry.size.size), f.frag_offset); - // } - - // break :blk rdr; - // }, - // else => unreachable, - // }; - - // var atomic = try Io.Dir.cwd().createFileAtomic(io, path, .{}); - // defer atomic.deinit(io); - - // try ext.extractAsync(alloc, io, atomic.file); - - // try atomic.link(io); - - var rdr: DataReader = switch (inode.data) { - .file => |f| blk: { - var rdr: DataReader = .init(alloc, data, decomp, block_size, f.blocks, f.size, f.block_start); - rdr.addCache(io, cache); - - if (f.frag_idx == 0xFFFFFFFF) break :blk rdr; - - const entry: Lookup.FragEntry = try frag_table.get(io, f.frag_idx); - if (entry.size.uncompressed) { - rdr.addFrag(data[entry.block_start..][0..entry.size.size], f.frag_offset); - } else { - rdr.addFrag(try cache.get(io, entry.block_start, entry.size.size), f.frag_offset); - } - - break :blk rdr; - }, - .ext_file => |f| blk: { - if (f.xattr_idx != 0xFFFFFFFF) ret.xattr_idx = f.xattr_idx; - - var rdr: DataReader = .init(alloc, data, decomp, block_size, f.blocks, f.size, f.block_start); - rdr.addCache(io, cache); - - if (f.frag_idx == 0xFFFFFFFF) break :blk rdr; - - const entry: Lookup.FragEntry = try frag_table.get(io, f.frag_idx); - if (entry.size.uncompressed) { - rdr.addFrag(data[entry.block_start..][0..entry.size.size], f.frag_offset); - } else { - rdr.addFrag(try cache.get(io, entry.block_start, entry.size.size), f.frag_offset); - } - - break :blk rdr; - }, - else => unreachable, - }; - - var atomic = try Io.Dir.cwd().createFileAtomic(io, path, .{}); - defer atomic.deinit(io); - - var writer = atomic.file.writer(io, &[0]u8{}); - _ = try rdr.interface.streamRemaining(&writer.interface); - try writer.flush(); - - try atomic.link(io); - - return ret; -} -fn extractSymlink(alloc: std.mem.Allocator, io: Io, inode: Inode, path: []const u8, origin: bool) Error!void { - defer if (!origin) { - inode.deinit(alloc); - alloc.free(path); - }; - - const target = switch (inode.data) { - .symlink => |s| s.target, - .ext_symlink => |s| s.target, - else => unreachable, - }; - try Io.Dir.cwd().symLink(io, target, path, .{}); -} -fn extractNod(alloc: std.mem.Allocator, inode: Inode, path: []const u8, origin: bool) Error!PathReturn { - errdefer if (!origin) - alloc.free(path); - - var ret: PathReturn = .{ - .hdr = inode.hdr, - .path = path, - }; - - var dev: u32 = 0; - var mode: u32 = undefined; - - const DT = std.posix.DT; - - switch (inode.data) { - .char_dev => |d| { - dev = d.device; - mode = DT.CHR; - }, - .block_dev => |d| { - dev = d.device; - mode = DT.BLK; - }, - .ext_char_dev => |d| { - if (d.xattr_idx != 0xFFFFFFFF) ret.xattr_idx = d.xattr_idx; - - dev = d.device; - mode = DT.CHR; - }, - .ext_block_dev => |d| { - if (d.xattr_idx != 0xFFFFFFFF) ret.xattr_idx = d.xattr_idx; - - dev = d.device; - mode = DT.BLK; - }, - .fifo => mode = DT.FIFO, - .socket => mode = DT.SOCK, - .ext_fifo => |i| { - if (i.xattr_idx != 0xFFFFFFFF) ret.xattr_idx = i.xattr_idx; - - mode = DT.FIFO; - }, - .ext_socket => |i| { - if (i.xattr_idx != 0xFFFFFFFF) ret.xattr_idx = i.xattr_idx; - - mode = DT.SOCK; - }, - else => unreachable, + for (xattr.kvs) |kv| { + const res = std.os.linux.fsetxattr(fil.handle, kv.key, kv.value.ptr, kv.value.len, 0); + if (res != 0) + return error.SetXattrError; + } + } + if (!options.ignore_permissions) { + try fil.setTimestamps(io, .{ + .modify_timestamp = .init(Io.Timestamp.fromNanoseconds(@as(i96, @intCast(inode.hdr.mod_time)) * std.time.ns_per_s)), + }); + try fil.setPermissions(io, @enumFromInt(inode.hdr.permissions)); + try fil.setOwner(io, try id_table.get(io, inode.hdr.uid_idx), try id_table.get(io, inode.hdr.gid_idx)); } - - const sentinel_path = try alloc.dupeSentinel(u8, path, 0); - defer alloc.free(sentinel_path); - - const res = std.os.linux.mknod(sentinel_path, mode, dev); - if (res != 0) - return error.MknodError; - - return ret; } // Types diff --git a/src/extract-single.zig b/src/extract-single.zig index c8b0d7e..b3fe1b4 100644 --- a/src/extract-single.zig +++ b/src/extract-single.zig @@ -183,6 +183,7 @@ fn extractFile( }, else => unreachable, }; + defer rdr.deinit(); var atomic = try Io.Dir.cwd().createFileAtomic(io, path, .{}); defer atomic.deinit(io);