| ... | ... | @@ -210,25 +210,38 @@ pub fn connect( |
| 210 | 210 | var connect_many_queue: Io.Queue(ConnectManyResult) = .init(&connect_many_buffer); |
| 211 | 211 | |
| 212 | 212 | var connect_many = io.async(connectMany, .{ host_name, io, port, &connect_many_queue, options }); |
| 213 | | defer connect_many.cancel(io); |
| 213 | var saw_end = false; |
| 214 | defer { |
| 215 | connect_many.cancel(io); |
| 216 | if (!saw_end) while (true) switch (connect_many_queue.getOneUncancelable(io)) { |
| 217 | .connection => |loser| if (loser) |s| s.closeConst(io) else |_| continue, |
| 218 | .end => break, |
| 219 | }; |
| 220 | } |
| 214 | 221 | |
| 215 | 222 | var aggregate_error: ConnectError = error.UnknownHostName; |
| 216 | 223 | |
| 217 | 224 | while (connect_many_queue.getOne(io)) |result| switch (result) { |
| 218 | 225 | .connection => |connection| if (connection) |stream| return stream else |err| switch (err) { |
| 219 | | error.SystemResources => |e| return e, |
| 220 | | error.OptionUnsupported => |e| return e, |
| 221 | | error.ProcessFdQuotaExceeded => |e| return e, |
| 222 | | error.SystemFdQuotaExceeded => |e| return e, |
| 223 | | error.Canceled => |e| return e, |
| 226 | error.SystemResources, |
| 227 | error.OptionUnsupported, |
| 228 | error.ProcessFdQuotaExceeded, |
| 229 | error.SystemFdQuotaExceeded, |
| 230 | error.Canceled, |
| 231 | => |e| return e, |
| 232 | |
| 224 | 233 | error.WouldBlock => return error.Unexpected, |
| 234 | |
| 225 | 235 | else => |e| aggregate_error = e, |
| 226 | 236 | }, |
| 227 | 237 | .end => |end| { |
| 238 | saw_end = true; |
| 228 | 239 | try end; |
| 229 | 240 | return aggregate_error; |
| 230 | 241 | }, |
| 231 | | } else |err| return err; |
| 242 | } else |err| switch (err) { |
| 243 | error.Canceled => |e| return e, |
| 244 | } |
| 232 | 245 | } |
| 233 | 246 | |
| 234 | 247 | pub const ConnectManyResult = union(enum) { |
| ... | ... | @@ -255,7 +268,6 @@ pub fn connectMany( |
| 255 | 268 | }); |
| 256 | 269 | |
| 257 | 270 | var group: Io.Group = .init; |
| 258 | | defer group.cancel(io); |
| 259 | 271 | |
| 260 | 272 | while (lookup_queue.getOne(io)) |dns_result| switch (dns_result) { |
| 261 | 273 | .address => |address| group.async(io, enqueueConnection, .{ address, io, results, options }), |
| ... | ... | @@ -266,7 +278,10 @@ pub fn connectMany( |
| 266 | 278 | return; |
| 267 | 279 | }, |
| 268 | 280 | } else |err| switch (err) { |
| 269 | | error.Canceled => |e| results.putOneUncancelable(io, .{ .end = e }), |
| 281 | error.Canceled => |e| { |
| 282 | group.cancel(io); |
| 283 | results.putOneUncancelable(io, .{ .end = e }); |
| 284 | }, |
| 270 | 285 | } |
| 271 | 286 | } |
| 272 | 287 | |