diff --git a/src/internal/subscriber.zig b/src/internal/subscriber.zig index 2b195d6..944f426 100644 --- a/src/internal/subscriber.zig +++ b/src/internal/subscriber.zig @@ -300,6 +300,37 @@ pub const Subscriber = struct { } } + /// resolve and connect trying addresses one at a time, instead of + /// HostName.connect's racing per-address attempts (Io.Group + loser + /// cancellation). a relay reconnecting to thousands of hosts wants the + /// smallest possible concurrency surface per connect, not per-connect + /// latency: sequential attempts keep the subscriber fiber the only + /// actor, with no group spawns and no cancel storms in the hot path. + fn connectSequential(io: Io, host_name: Io.net.HostName) !Io.net.Stream { + var canonical_name_buffer: [Io.net.HostName.max_len]u8 = undefined; + var lookup_buffer: [32]Io.net.HostName.LookupResult = undefined; + var lookup_queue: Io.Queue(Io.net.HostName.LookupResult) = .init(&lookup_buffer); + try host_name.lookup(io, &lookup_queue, .{ + .port = 443, + .canonical_name_buffer = &canonical_name_buffer, + }); + var last_err: ?anyerror = null; + while (lookup_queue.getOne(io)) |result| switch (result) { + .address => |address| { + const stream = address.connect(io, .{ .mode = .stream }) catch |e| { + last_err = e; + continue; + }; + return stream; + }, + .canonical_name => continue, + } else |err| switch (err) { + error.Closed => {}, + else => |e| return e, + } + return if (last_err) |e| e else error.ConnectionRefused; + } + fn connectAndRead(self: *Subscriber) !void { var path_buf: [256]u8 = undefined; var w: std.Io.Writer = .fixed(&path_buf); @@ -324,7 +355,7 @@ pub const Subscriber = struct { _ = self.bc.stats.subscriber_disconnect_dns_connect.fetchAdd(1, .monotonic); return e; }; - const net_stream = host_name.connect(self.io, 443, .{ .mode = .stream }) catch |e| { + const net_stream = connectSequential(self.io, host_name) catch |e| { _ = self.bc.stats.subscriber_disconnect_dns_connect.fetchAdd(1, .monotonic); return e; };