Re-working multi-threaded extraction

This commit is contained in:
Caleb Gardner
2026-06-20 06:34:33 -05:00
parent aa6b61fc9f
commit 86c9829a23
3 changed files with 139 additions and 339 deletions
+17 -3
View File
@@ -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;
+121 -336
View File
@@ -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
+1
View File
@@ -183,6 +183,7 @@ fn extractFile(
},
else => unreachable,
};
defer rdr.deinit();
var atomic = try Io.Dir.cwd().createFileAtomic(io, path, .{});
defer atomic.deinit(io);