Finished (probably not actually) mutli-threaded extraction

This commit is contained in:
Caleb Gardner
2026-07-11 00:47:39 -05:00
parent 7058f737c2
commit 28f7d1fe28
5 changed files with 210 additions and 52 deletions
+9 -8
View File
@@ -2,7 +2,8 @@ const std = @import("std");
const Io = std.Io; const Io = std.Io;
const Decomp = @import("../decomp.zig"); const Decomp = @import("../decomp.zig");
const Finish = @import("../extract-multi.zig").Finish; const FileFinish = @import("../extract-multi.zig").FileFinish;
const ExtractError = @import("../extract-multi.zig").Error;
const DataBlock = @import("../inode.zig").DataBlock; const DataBlock = @import("../inode.zig").DataBlock;
const Cache = @import("../util/cache.zig"); const Cache = @import("../util/cache.zig");
@@ -42,16 +43,16 @@ pub fn addCache(self: *Extractor, cache: *Cache) void {
self.cache = cache; self.cache = cache;
} }
pub fn extractAsync(self: Extractor, alloc: std.mem.Allocator, io: Io, file: Io.File, group: *Io.Group, err: *?Error, finish: *Finish) Error!void { pub fn extractAsync(self: Extractor, alloc: std.mem.Allocator, io: Io, group: *Io.Group, err: *?ExtractError, finish: *FileFinish) void {
if (self.size == 0) return; if (self.size == 0) return;
var read_offset: u64 = self.start; var read_offset: u64 = self.start;
for (0.., self.blocks) |i, block| { for (0.., self.blocks) |i, block| {
group.async(io, blockThread, .{ self, alloc, io, file, read_offset, @truncate(i), err, finish }); group.async(io, blockThread, .{ self, alloc, io, finish.file.file, read_offset, @truncate(i), err, finish });
read_offset += block.size; read_offset += block.size;
} }
if (self.frag_data != null) if (self.frag_data != null)
group.async(io, fragThread, .{ self, io, file, err, finish }); group.async(io, fragThread, .{ self, io, finish.file.file, err, finish });
} }
fn blockThread( fn blockThread(
@@ -61,8 +62,8 @@ fn blockThread(
file: Io.File, file: Io.File,
read_offset: u64, read_offset: u64,
block_idx: u32, block_idx: u32,
err: *?Error, err: *?ExtractError,
finish: *Finish, finish: *FileFinish,
) error{Canceled}!void { ) error{Canceled}!void {
const size = if (self.frag_data == null and block_idx == self.blocks.len - 1) const size = if (self.frag_data == null and block_idx == self.blocks.len - 1)
self.size % self.block_size self.size % self.block_size
@@ -135,8 +136,8 @@ fn fragThread(
self: Extractor, self: Extractor,
io: Io, io: Io,
file: Io.File, file: Io.File,
err: *?Error, err: *?ExtractError,
finish: *Finish, finish: *FileFinish,
) error{Canceled}!void { ) error{Canceled}!void {
const size = self.size % self.block_size; const size = self.size % self.block_size;
+195 -39
View File
@@ -2,7 +2,6 @@ const std = @import("std");
const Io = std.Io; const Io = std.Io;
const Atomic = std.atomic.Value; const Atomic = std.atomic.Value;
const Archive = @import("archive.zig");
const DataExtractor = @import("data/extractor.zig"); const DataExtractor = @import("data/extractor.zig");
const Decomp = @import("decomp.zig"); const Decomp = @import("decomp.zig");
const Directory = @import("directory.zig"); const Directory = @import("directory.zig");
@@ -10,46 +9,57 @@ const ExtractionOption = @import("options.zig");
const Inode = @import("inode.zig"); const Inode = @import("inode.zig");
const Lookup = @import("lookup.zig"); const Lookup = @import("lookup.zig");
const MetadataReader = @import("meta_rdr.zig"); const MetadataReader = @import("meta_rdr.zig");
const Superblock = @import("archive.zig").Superblock;
const Cache = @import("util/cache.zig");
const XattrTable = @import("xattr.zig"); const XattrTable = @import("xattr.zig");
pub fn extract(alloc: std.mem.Allocator, io: Io, arc: Archive, options: ExtractionOption) !void { pub fn extract(alloc: std.mem.Allocator, io: Io, super: Superblock, data: []u8, decomp: Decomp.Fn, inode: Inode, filepath: []const u8, options: ExtractionOption) !void {
var common: Common = .init(alloc, arc, options); const path = std.mem.trim(u8, filepath, "/");
var common: Common = .init(alloc, super, data, decomp, options);
defer common.deinit(); defer common.deinit();
common.group.async(io, extractAsync, .{ &common, io, inode, path, null });
try common.group.await(io);
if (common.err != null)
return common.err.?;
} }
fn extractAsync(common: *Common, io: Io, inode: Inode, path: []const u8, parent: ?*Finish) error{Canceled}!void { fn extractAsync(common: *Common, io: Io, inode: Inode, path: []const u8, parent: ?*Parent) error{Canceled}!void {
switch (inode.hdr.type) { switch (inode.hdr.type) {
.dir, .ext_dir => extractDir(common, io, inode, path, parent) catch |err| switch (err) { .dir, .ext_dir => extractDir(common, io, inode, path, parent) catch |err| switch (err) {
error.Canceled => { error.Canceled => {
io.recancel(io); io.recancel();
return error.Canceled; return error.Canceled;
}, },
else => common.err = err, else => common.err = err,
}, },
.file, .ext_file => extractReg(common, io, inode, path, parent) catch |err| switch (err) { .file, .ext_file => extractReg(common, io, inode, path, parent) catch |err| switch (err) {
error.Canceled => { error.Canceled => {
io.recancel(io); io.recancel();
return error.Canceled; return error.Canceled;
}, },
else => common.err = err, else => common.err = err,
}, },
.symlink, .ext_symlink => extractReg(common, io, inode, path, parent) catch |err| switch (err) { .symlink, .ext_symlink => extractReg(common, io, inode, path, parent) catch |err| switch (err) {
error.Canceled => { error.Canceled => {
io.recancel(io); io.recancel();
return error.Canceled; return error.Canceled;
}, },
else => common.err = err, else => common.err = err,
}, },
else => extractNod(common, io, inode, path, parent) catch |err| switch (err) { else => extractNod(common, io, inode, path, parent) catch |err| switch (err) {
error.Canceled => { error.Canceled => {
io.recancel(io); io.recancel();
return error.Canceled; return error.Canceled;
}, },
else => common.err = err, else => common.err = err,
}, },
} }
} }
fn extractDir(common: *Common, io: Io, inode: Inode, path: []const u8, parent: ?*Finish) Error!void { fn extractDir(common: *Common, io: Io, inode: Inode, path: []const u8, parent: ?*Parent) Error!void {
var xattr_idx: u32 = 0xFFFFFFFF; var xattr_idx: u32 = 0xFFFFFFFF;
const dir: Directory = switch (inode.data) { const dir: Directory = switch (inode.data) {
@@ -74,15 +84,15 @@ fn extractDir(common: *Common, io: Io, inode: Inode, path: []const u8, parent: ?
try Io.Dir.cwd().createDirPath(io, path); try Io.Dir.cwd().createDirPath(io, path);
if (dir.entries.len == 0) { if (dir.entries.len == 0) {
try setMetadataPath(io, inode.hdr, path, xattr_idx); try setMetadataPath(common, io, inode.hdr, path, xattr_idx);
if (parent != null) { if (parent != null) {
try parent.?.finish(); try parent.?.finish(io);
common.alloc.free(path); common.alloc.free(path);
} }
return; return;
} }
const cur_dir: *Finish = try .init(common, dir.entries.len, path, xattr_idx, parent); const cur_dir: *Parent = try .init(common, dir.entries.len, inode.hdr, path, xattr_idx, parent);
for (dir.entries) |entry| { for (dir.entries) |entry| {
var new_inode: Inode = try .initEntry(common.alloc, common.data, common.decomp, common.inode_start, common.block_size, entry); var new_inode: Inode = try .initEntry(common.alloc, common.data, common.decomp, common.inode_start, common.block_size, entry);
@@ -95,7 +105,7 @@ fn extractDir(common: *Common, io: Io, inode: Inode, path: []const u8, parent: ?
common.group.async(io, extractAsync, .{ common, io, new_inode, new_path, cur_dir }); common.group.async(io, extractAsync, .{ common, io, new_inode, new_path, cur_dir });
} }
} }
fn extractReg(common: *Common, io: Io, inode: Inode, path: []const u8, parent: ?*Finish) Error!void { fn extractReg(common: *Common, io: Io, inode: Inode, path: []const u8, parent: ?*Parent) Error!void {
var blocks: usize = 0; var blocks: usize = 0;
var xattr_idx: u32 = 0xFFFFFFFF; var xattr_idx: u32 = 0xFFFFFFFF;
@@ -104,22 +114,122 @@ fn extractReg(common: *Common, io: Io, inode: Inode, path: []const u8, parent: ?
var ext: DataExtractor = .init(common.data, common.decomp, common.block_size, f.blocks, f.block_start, f.size); var ext: DataExtractor = .init(common.data, common.decomp, common.block_size, f.blocks, f.block_start, f.size);
blocks = f.blocks.len; blocks = f.blocks.len;
if (f.frag_idx != 0xFFFFFFFF) {} if (f.frag_idx != 0xFFFFFFFF) {
blocks += 1;
const frag: Lookup.FragEntry = try common.frag_table.get(io, f.frag_idx);
if (frag.size.uncompressed) {
ext.addFrag(common.data[frag.block_start..][0 .. f.size % common.block_size], f.frag_offset);
} else {
const frag_block = try common.cache.get(io, frag.block_start, frag.size.size);
ext.addFrag(frag_block, f.frag_offset);
}
}
break :blk ext;
}, },
.ext_file => |f| blk: { .ext_file => |f| blk: {
xattr_idx = f.xattr_idx; xattr_idx = f.xattr_idx;
var ext: DataExtractor = .init(common.data, common.decomp, common.block_size, f.blocks, f.block_start, f.size);
blocks = f.blocks.len;
if (f.frag_idx != 0xFFFFFFFF) {
blocks += 1;
const frag: Lookup.FragEntry = try common.frag_table.get(io, f.frag_idx);
if (frag.size.uncompressed) {
ext.addFrag(common.data[frag.block_start..][0 .. f.size % common.block_size], f.frag_offset);
} else {
const frag_block = try common.cache.get(io, frag.block_start, frag.size.size);
ext.addFrag(frag_block, f.frag_offset);
}
}
break :blk ext;
}, },
else => unreachable, else => unreachable,
}; };
var fin: *Finish = try .init(common, blocks, inode.hdr, path, xattr_idx, parent); const fin: *FileFinish = try .init(common, io, blocks, inode.hdr, path, xattr_idx, parent);
ext.extractAsync(common.alloc, io, &common.group, @ptrCast(&common.err), fin);
}
fn extractSymlink(common: *Common, io: Io, inode: Inode, path: []const u8, parent: ?*Parent) Error!void {
defer if (parent != null) {
inode.deinit(common.alloc);
common.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, .{});
if (parent != null)
try parent.?.finish(io);
}
fn extractNod(common: *Common, io: Io, inode: Inode, path: []const u8, parent: ?*Parent) Error!void {
defer if (parent != null) {};
var xattr_idx: u32 = 0xFFFFFFFF;
var device: u32 = 0;
var mode: u32 = undefined;
const DT = std.os.linux.DT;
switch (inode.data) {
.block_dev => |d| {
mode = DT.BLK;
device = d.device;
},
.ext_block_dev => |d| {
mode = DT.BLK;
device = d.device;
xattr_idx = d.xattr_idx;
},
.char_dev => |d| {
mode = DT.CHR;
device = d.device;
},
.ext_char_dev => |d| {
mode = DT.CHR;
device = d.device;
xattr_idx = d.xattr_idx;
},
.fifo => {
mode = DT.FIFO;
},
.ext_fifo => |f| {
mode = DT.FIFO;
xattr_idx = f.xattr_idx;
},
.socket => {
mode = DT.SOCK;
},
.ext_socket => |s| {
mode = DT.SOCK;
xattr_idx = s.xattr_idx;
},
else => unreachable,
}
const sentinel_path = try common.alloc.dupeSentinel(u8, path, 0);
defer common.alloc.free(sentinel_path);
const res = std.os.linux.mknod(sentinel_path, mode, device);
if (res != 0)
return Error.Mknod;
try setMetadataPath(common, io, inode.hdr, path, xattr_idx);
if (parent != null)
try parent.?.finish(io);
} }
fn setMetadataPath(common: *Common, io: Io, hdr: Inode.Header, path: []const u8, xattr_idx: u32) Error!void { fn setMetadataPath(common: *Common, io: Io, hdr: Inode.Header, path: []const u8, xattr_idx: u32) Error!void {
if (common.options.ignore_permissions and (common.options.ignore_xattr or xattr_idx == 0xFFFFFFFF)) return; if (common.options.ignore_permissions and (common.options.ignore_xattr or xattr_idx == 0xFFFFFFFF)) return;
var file = try Io.Dir.cwd().openFile(io, path, .{}); var file = try Io.Dir.cwd().openFile(io, path, .{});
defer file.deinit(); defer file.close(io);
return setMetadata(common, io, hdr, file, xattr_idx); return setMetadata(common, io, hdr, file, xattr_idx);
} }
@@ -131,7 +241,7 @@ fn setMetadata(common: *Common, io: Io, hdr: Inode.Header, file: Io.File, xattr_
defer xattr.deinit(common.alloc); defer xattr.deinit(common.alloc);
for (xattr.kvs) |kv| { for (xattr.kvs) |kv| {
const res = std.os.linux.fsetxattr(file.handle, kv.key, kv.value, kv.value.len, 0); const res = std.os.linux.fsetxattr(file.handle, kv.key, kv.value.ptr, kv.value.len, 0);
if (res != 0) if (res != 0)
return Error.SetXattr; return Error.SetXattr;
} }
@@ -143,14 +253,15 @@ fn setMetadata(common: *Common, io: Io, hdr: Inode.Header, file: Io.File, xattr_
.modify_timestamp = .init(Io.Timestamp.fromNanoseconds(hdr.mod_time * std.time.ns_per_s)), .modify_timestamp = .init(Io.Timestamp.fromNanoseconds(hdr.mod_time * std.time.ns_per_s)),
}); });
try file.setPermissions(io, @enumFromInt(hdr.permissions)); try file.setPermissions(io, @enumFromInt(hdr.permissions));
try file.setOwner(io, try common.id_table.get(hdr.uid_idx), try common.id_table.get(hdr.gid_idx)); try file.setOwner(io, try common.id_table.get(io, hdr.uid_idx), try common.id_table.get(io, hdr.gid_idx));
} }
// Types // Types
const Error = error{ Mknod, SetXattr } || std.mem.Allocator.Error; pub const Error = error{ Mknod, SetXattr } || std.mem.Allocator.Error || Io.Cancelable || Io.Reader.Error || Io.Dir.CreateDirPathError || Cache.Error ||
Io.File.SetPermissionsError || DataExtractor.Error;
pub const Finish = struct { pub const Parent = struct {
common: *Common, common: *Common,
atomic: Atomic(usize), atomic: Atomic(usize),
@@ -158,10 +269,10 @@ pub const Finish = struct {
path: []const u8, path: []const u8,
xattr_idx: u32, xattr_idx: u32,
parent: ?*Finish, parent: ?*Parent,
fn init(common: *Common, len: usize, hdr: Inode.Header, path: []const u8, xattr_idx: u32, parent: ?*Finish) !*Finish { fn init(common: *Common, len: usize, hdr: Inode.Header, path: []const u8, xattr_idx: u32, parent: ?*Parent) !*Parent {
const out = try common.dirAlloc().create(Finish); const out = try common.arenaAlloc().create(Parent);
out.* = .{ out.* = .{
.common = common, .common = common,
@@ -176,8 +287,8 @@ pub const Finish = struct {
return out; return out;
} }
pub fn finish(self: *Finish, io: Io) Error!void { pub fn finish(self: *Parent, io: Io) Error!void {
if (self.atomic.fetchSub(1, .unordered) >= 0) return; if (self.atomic.fetchSub(1, .release) >= 0) return;
try setMetadataPath(self.common, io, self.hdr, self.path, self.xattr_idx); try setMetadataPath(self.common, io, self.hdr, self.path, self.xattr_idx);
@@ -188,10 +299,52 @@ pub const Finish = struct {
} }
} }
}; };
pub const FileFinish = struct {
common: *Common,
atomic: Atomic(usize),
hdr: Inode.Header,
file: Io.File.Atomic,
xattr_idx: u32,
parent: ?*Parent,
fn init(common: *Common, io: Io, len: usize, hdr: Inode.Header, path: []const u8, xattr_idx: u32, parent: ?*Parent) !*FileFinish {
const out = try common.arenaAlloc().create(FileFinish);
out.* = .{
.common = common,
.atomic = .init(len),
.hdr = hdr,
.file = try Io.Dir.cwd().createFileAtomic(io, path, .{}),
.xattr_idx = xattr_idx,
.parent = parent,
};
return out;
}
pub fn finish(self: *FileFinish, io: Io) Error!void {
if (self.atomic.fetchSub(1, .release) >= 0) return;
defer self.file.deinit(io);
try setMetadata(self.common, io, self.hdr, self.file.file, self.xattr_idx);
try self.file.link(io);
if (self.parent != null) {
if (self.hdr.type != .dir and self.hdr.type != .ext_dir)
self.common.alloc.free(self.path);
try self.parent.?.finish(io);
}
}
};
const Common = struct { const Common = struct {
alloc: std.mem.Allocator, alloc: std.mem.Allocator,
dir_alloc: std.heap.ArenaAllocator, arena: std.heap.ArenaAllocator,
group: Io.Group = .init, group: Io.Group = .init,
err: ?Error = null, err: ?Error = null,
@@ -206,37 +359,40 @@ const Common = struct {
id_table: Lookup.Table(u16), id_table: Lookup.Table(u16),
frag_table: Lookup.Table(Lookup.FragEntry), frag_table: Lookup.Table(Lookup.FragEntry),
xattr_table: XattrTable, xattr_table: XattrTable,
cache: Cache,
options: ExtractionOption, options: ExtractionOption,
fn init(alloc: std.mem.Allocator, arc: Archive, options: ExtractionOption) Common { fn init(alloc: std.mem.Allocator, super: Superblock, data: []u8, decomp: Decomp.Fn, options: ExtractionOption) Common {
return .{ return .{
.alloc = alloc, .alloc = alloc,
.dir_alloc = .init(alloc), .arena = .init(alloc),
.data = arc.map.memory, .data = data,
.decomp = arc.decomp, .decomp = decomp,
.inode_start = arc.super.inode_start, .inode_start = super.inode_start,
.dir_start = arc.super.dir_start, .dir_start = super.dir_start,
.block_size = arc.super.block_size, .block_size = super.block_size,
.id_table = .init(alloc, arc.map.memory, arc.decomp, arc.super.id_start, arc.super.id_count), .id_table = .init(alloc, data, decomp, super.id_start, super.id_count),
.frag_table = .init(alloc, arc.map.memory, arc.decomp, arc.super.frag_start, arc.super.frag_count), .frag_table = .init(alloc, data, decomp, super.frag_start, super.frag_count),
.xattr_table = .init(alloc, arc.map.memory, arc.decomp, arc.super.xattr_start), .xattr_table = .init(alloc, data, decomp, super.xattr_start),
.cache = .init(alloc, data, decomp),
.options = options, .options = options,
}; };
} }
fn deinit(self: *Common) void { fn deinit(self: *Common) void {
self.dir_alloc.deinit(); self.arena.deinit();
self.id_table.deinit(); self.id_table.deinit();
self.frag_table.deinit(); self.frag_table.deinit();
self.xattr_table.deinit(); self.xattr_table.deinit();
self.cache.deinit();
} }
fn dirAlloc(self: *Common) std.mem.Allocator { fn arenaAlloc(self: *Common) std.mem.Allocator {
return self.dir_alloc.allocator(); return self.arena.allocator();
} }
}; };
+2 -1
View File
@@ -1,6 +1,8 @@
const std = @import("std"); const std = @import("std");
const Io = std.Io; const Io = std.Io;
const Archive = @import("archive.zig");
const Superblock = Archive.Superblock;
const DataReader = @import("data/reader.zig"); const DataReader = @import("data/reader.zig");
const Decomp = @import("decomp.zig"); const Decomp = @import("decomp.zig");
const Directory = @import("directory.zig"); const Directory = @import("directory.zig");
@@ -8,7 +10,6 @@ const ExtractionOptions = @import("options.zig");
const Inode = @import("inode.zig"); const Inode = @import("inode.zig");
const Lookup = @import("lookup.zig"); const Lookup = @import("lookup.zig");
const MetadaReader = @import("meta_rdr.zig"); const MetadaReader = @import("meta_rdr.zig");
const Superblock = @import("archive.zig").Superblock;
const Cache = @import("util/cache.zig"); const Cache = @import("util/cache.zig");
const XattrTable = @import("xattr.zig"); const XattrTable = @import("xattr.zig");
+3 -3
View File
@@ -1,12 +1,12 @@
const std = @import("std"); const std = @import("std");
const Io = std.Io; const Io = std.Io;
const Superblock = @import("archive.zig").Superblock;
const Decomp = @import("decomp.zig"); const Decomp = @import("decomp.zig");
const Inode = @import("inode.zig");
const ExtractionOptions = @import("options.zig"); const ExtractionOptions = @import("options.zig");
const Single = @import("extract-single.zig"); const Inode = @import("inode.zig");
const Multi = @import("extract-multi.zig"); const Multi = @import("extract-multi.zig");
const Single = @import("extract-single.zig");
const Superblock = @import("archive.zig").Superblock;
pub fn extract( pub fn extract(
alloc: std.mem.Allocator, alloc: std.mem.Allocator,
+1 -1
View File
@@ -38,7 +38,7 @@ pub fn ProtectedMap(comptime K: anytype, comptime T: anytype, comptime create_fn
if (value != null) { if (value != null) {
if (!value.?.filled.isSet()) { if (!value.?.filled.isSet()) {
self.mut.unlockShared(io); self.mut.unlockShared(io);
defer self.mut.lockShared(io); defer self.mut.lockSharedUncancelable(io);
try value.?.filled.wait(io); try value.?.filled.wait(io);
} }