| ... | @@ -10,6 +10,9 @@ const Loop = event.Loop; | ... | @@ -10,6 +10,9 @@ const Loop = event.Loop; |
| 10 | const fd_t = os.fd_t; | 10 | const fd_t = os.fd_t; |
| 11 | const File = std.fs.File; | 11 | const File = std.fs.File; |
| 12 | | 12 | |
| | 13 | const global_event_loop = Loop.instance orelse |
| | 14 | @compileError("std.event.fs currently only works with event-based I/O"); |
| | 15 | |
| 13 | pub const RequestNode = std.atomic.Queue(Request).Node; | 16 | pub const RequestNode = std.atomic.Queue(Request).Node; |
| 14 | | 17 | |
| 15 | pub const Request = struct { | 18 | pub const Request = struct { |
| ... | @@ -86,7 +89,7 @@ pub const Request = struct { | ... | @@ -86,7 +89,7 @@ pub const Request = struct { |
| 86 | pub const PWriteVError = error{OutOfMemory} || File.WriteError; | 89 | pub const PWriteVError = error{OutOfMemory} || File.WriteError; |
| 87 | | 90 | |
| 88 | /// data - just the inner references - must live until pwritev frame completes. | 91 | /// data - just the inner references - must live until pwritev frame completes. |
| 89 | pub fn pwritev(loop: *Loop, fd: fd_t, data: []const []const u8, offset: usize) PWriteVError!void { | 92 | pub fn pwritev(allocator: *Allocator, fd: fd_t, data: []const []const u8, offset: usize) PWriteVError!void { |
| 90 | switch (builtin.os) { | 93 | switch (builtin.os) { |
| 91 | .macosx, | 94 | .macosx, |
| 92 | .linux, | 95 | .linux, |
| ... | @@ -94,8 +97,8 @@ pub fn pwritev(loop: *Loop, fd: fd_t, data: []const []const u8, offset: usize) P | ... | @@ -94,8 +97,8 @@ pub fn pwritev(loop: *Loop, fd: fd_t, data: []const []const u8, offset: usize) P |
| 94 | .netbsd, | 97 | .netbsd, |
| 95 | .dragonfly, | 98 | .dragonfly, |
| 96 | => { | 99 | => { |
| 97 | const iovecs = try loop.allocator.alloc(os.iovec_const, data.len); | 100 | const iovecs = try allocator.alloc(os.iovec_const, data.len); |
| 98 | defer loop.allocator.free(iovecs); | 101 | defer allocator.free(iovecs); |
| 99 | | 102 | |
| 100 | for (data) |buf, i| { | 103 | for (data) |buf, i| { |
| 101 | iovecs[i] = os.iovec_const{ | 104 | iovecs[i] = os.iovec_const{ |
| ... | @@ -104,31 +107,31 @@ pub fn pwritev(loop: *Loop, fd: fd_t, data: []const []const u8, offset: usize) P | ... | @@ -104,31 +107,31 @@ pub fn pwritev(loop: *Loop, fd: fd_t, data: []const []const u8, offset: usize) P |
| 104 | }; | 107 | }; |
| 105 | } | 108 | } |
| 106 | | 109 | |
| 107 | return pwritevPosix(loop, fd, iovecs, offset); | 110 | return pwritevPosix(fd, iovecs, offset); |
| 108 | }, | 111 | }, |
| 109 | .windows => { | 112 | .windows => { |
| 110 | const data_copy = try std.mem.dupe(loop.allocator, []const u8, data); | 113 | const data_copy = try std.mem.dupe(allocator, []const u8, data); |
| 111 | defer loop.allocator.free(data_copy); | 114 | defer allocator.free(data_copy); |
| 112 | return pwritevWindows(loop, fd, data, offset); | 115 | return pwritevWindows(fd, data, offset); |
| 113 | }, | 116 | }, |
| 114 | else => @compileError("Unsupported OS"), | 117 | else => @compileError("Unsupported OS"), |
| 115 | } | 118 | } |
| 116 | } | 119 | } |
| 117 | | 120 | |
| 118 | /// data must outlive the returned frame | 121 | /// data must outlive the returned frame |
| 119 | pub fn pwritevWindows(loop: *Loop, fd: fd_t, data: []const []const u8, offset: usize) os.WindowsWriteError!void { | 122 | pub fn pwritevWindows(fd: fd_t, data: []const []const u8, offset: usize) os.WindowsWriteError!void { |
| 120 | if (data.len == 0) return; | 123 | if (data.len == 0) return; |
| 121 | if (data.len == 1) return pwriteWindows(loop, fd, data[0], offset); | 124 | if (data.len == 1) return pwriteWindows(fd, data[0], offset); |
| 122 | | 125 | |
| 123 | // TODO do these in parallel | 126 | // TODO do these in parallel |
| 124 | var off = offset; | 127 | var off = offset; |
| 125 | for (data) |buf| { | 128 | for (data) |buf| { |
| 126 | try pwriteWindows(loop, fd, buf, off); | 129 | try pwriteWindows(fd, buf, off); |
| 127 | off += buf.len; | 130 | off += buf.len; |
| 128 | } | 131 | } |
| 129 | } | 132 | } |
| 130 | | 133 | |
| 131 | pub fn pwriteWindows(loop: *Loop, fd: fd_t, data: []const u8, offset: u64) os.WindowsWriteError!void { | 134 | pub fn pwriteWindows(fd: fd_t, data: []const u8, offset: u64) os.WindowsWriteError!void { |
| 132 | var resume_node = Loop.ResumeNode.Basic{ | 135 | var resume_node = Loop.ResumeNode.Basic{ |
| 133 | .base = Loop.ResumeNode{ | 136 | .base = Loop.ResumeNode{ |
| 134 | .id = Loop.ResumeNode.Id.Basic, | 137 | .id = Loop.ResumeNode.Id.Basic, |
| ... | @@ -143,9 +146,9 @@ pub fn pwriteWindows(loop: *Loop, fd: fd_t, data: []const u8, offset: u64) os.Wi | ... | @@ -143,9 +146,9 @@ pub fn pwriteWindows(loop: *Loop, fd: fd_t, data: []const u8, offset: u64) os.Wi |
| 143 | }, | 146 | }, |
| 144 | }; | 147 | }; |
| 145 | // TODO only call create io completion port once per fd | 148 | // TODO only call create io completion port once per fd |
| 146 | _ = windows.CreateIoCompletionPort(fd, loop.os_data.io_port, undefined, undefined); | 149 | _ = windows.CreateIoCompletionPort(fd, global_event_loop.os_data.io_port, undefined, undefined); |
| 147 | loop.beginOneEvent(); | 150 | global_event_loop.beginOneEvent(); |
| 148 | errdefer loop.finishOneEvent(); | 151 | errdefer global_event_loop.finishOneEvent(); |
| 149 | | 152 | |
| 150 | errdefer { | 153 | errdefer { |
| 151 | _ = windows.kernel32.CancelIoEx(fd, &resume_node.base.overlapped); | 154 | _ = windows.kernel32.CancelIoEx(fd, &resume_node.base.overlapped); |
| ... | @@ -168,12 +171,7 @@ pub fn pwriteWindows(loop: *Loop, fd: fd_t, data: []const u8, offset: u64) os.Wi | ... | @@ -168,12 +171,7 @@ pub fn pwriteWindows(loop: *Loop, fd: fd_t, data: []const u8, offset: u64) os.Wi |
| 168 | } | 171 | } |
| 169 | | 172 | |
| 170 | /// iovecs must live until pwritev frame completes. | 173 | /// iovecs must live until pwritev frame completes. |
| 171 | pub fn pwritevPosix( | 174 | pub fn pwritevPosix(fd: fd_t, iovecs: []const os.iovec_const, offset: usize) os.WriteError!void { |
| 172 | loop: *Loop, | | |
| 173 | fd: fd_t, | | |
| 174 | iovecs: []const os.iovec_const, | | |
| 175 | offset: usize, | | |
| 176 | ) os.WriteError!void { | | |
| 177 | var req_node = RequestNode{ | 175 | var req_node = RequestNode{ |
| 178 | .prev = null, | 176 | .prev = null, |
| 179 | .next = null, | 177 | .next = null, |
| ... | @@ -196,21 +194,17 @@ pub fn pwritevPosix( | ... | @@ -196,21 +194,17 @@ pub fn pwritevPosix( |
| 196 | }, | 194 | }, |
| 197 | }; | 195 | }; |
| 198 | | 196 | |
| 199 | errdefer loop.posixFsCancel(&req_node); | 197 | errdefer global_event_loop.posixFsCancel(&req_node); |
| 200 | | 198 | |
| 201 | suspend { | 199 | suspend { |
| 202 | loop.posixFsRequest(&req_node); | 200 | global_event_loop.posixFsRequest(&req_node); |
| 203 | } | 201 | } |
| 204 | | 202 | |
| 205 | return req_node.data.msg.PWriteV.result; | 203 | return req_node.data.msg.PWriteV.result; |
| 206 | } | 204 | } |
| 207 | | 205 | |
| 208 | /// iovecs must live until pwritev frame completes. | 206 | /// iovecs must live until pwritev frame completes. |
| 209 | pub fn writevPosix( | 207 | pub fn writevPosix(fd: fd_t, iovecs: []const os.iovec_const) os.WriteError!void { |
| 210 | loop: *Loop, | | |
| 211 | fd: fd_t, | | |
| 212 | iovecs: []const os.iovec_const, | | |
| 213 | ) os.WriteError!void { | | |
| 214 | var req_node = RequestNode{ | 208 | var req_node = RequestNode{ |
| 215 | .prev = null, | 209 | .prev = null, |
| 216 | .next = null, | 210 | .next = null, |
| ... | @@ -233,7 +227,7 @@ pub fn writevPosix( | ... | @@ -233,7 +227,7 @@ pub fn writevPosix( |
| 233 | }; | 227 | }; |
| 234 | | 228 | |
| 235 | suspend { | 229 | suspend { |
| 236 | loop.posixFsRequest(&req_node); | 230 | global_event_loop.posixFsRequest(&req_node); |
| 237 | } | 231 | } |
| 238 | | 232 | |
| 239 | return req_node.data.msg.WriteV.result; | 233 | return req_node.data.msg.WriteV.result; |
| ... | @@ -242,7 +236,7 @@ pub fn writevPosix( | ... | @@ -242,7 +236,7 @@ pub fn writevPosix( |
| 242 | pub const PReadVError = error{OutOfMemory} || File.ReadError; | 236 | pub const PReadVError = error{OutOfMemory} || File.ReadError; |
| 243 | | 237 | |
| 244 | /// data - just the inner references - must live until preadv frame completes. | 238 | /// data - just the inner references - must live until preadv frame completes. |
| 245 | pub fn preadv(loop: *Loop, fd: fd_t, data: []const []u8, offset: usize) PReadVError!usize { | 239 | pub fn preadv(allocator: *Allocator, fd: fd_t, data: []const []u8, offset: usize) PReadVError!usize { |
| 246 | assert(data.len != 0); | 240 | assert(data.len != 0); |
| 247 | switch (builtin.os) { | 241 | switch (builtin.os) { |
| 248 | .macosx, | 242 | .macosx, |
| ... | @@ -251,8 +245,8 @@ pub fn preadv(loop: *Loop, fd: fd_t, data: []const []u8, offset: usize) PReadVEr | ... | @@ -251,8 +245,8 @@ pub fn preadv(loop: *Loop, fd: fd_t, data: []const []u8, offset: usize) PReadVEr |
| 251 | .netbsd, | 245 | .netbsd, |
| 252 | .dragonfly, | 246 | .dragonfly, |
| 253 | => { | 247 | => { |
| 254 | const iovecs = try loop.allocator.alloc(os.iovec, data.len); | 248 | const iovecs = try allocator.alloc(os.iovec, data.len); |
| 255 | defer loop.allocator.free(iovecs); | 249 | defer allocator.free(iovecs); |
| 256 | | 250 | |
| 257 | for (data) |buf, i| { | 251 | for (data) |buf, i| { |
| 258 | iovecs[i] = os.iovec{ | 252 | iovecs[i] = os.iovec{ |
| ... | @@ -261,21 +255,21 @@ pub fn preadv(loop: *Loop, fd: fd_t, data: []const []u8, offset: usize) PReadVEr | ... | @@ -261,21 +255,21 @@ pub fn preadv(loop: *Loop, fd: fd_t, data: []const []u8, offset: usize) PReadVEr |
| 261 | }; | 255 | }; |
| 262 | } | 256 | } |
| 263 | | 257 | |
| 264 | return preadvPosix(loop, fd, iovecs, offset); | 258 | return preadvPosix(fd, iovecs, offset); |
| 265 | }, | 259 | }, |
| 266 | .windows => { | 260 | .windows => { |
| 267 | const data_copy = try std.mem.dupe(loop.allocator, []u8, data); | 261 | const data_copy = try std.mem.dupe(allocator, []u8, data); |
| 268 | defer loop.allocator.free(data_copy); | 262 | defer allocator.free(data_copy); |
| 269 | return preadvWindows(loop, fd, data_copy, offset); | 263 | return preadvWindows(fd, data_copy, offset); |
| 270 | }, | 264 | }, |
| 271 | else => @compileError("Unsupported OS"), | 265 | else => @compileError("Unsupported OS"), |
| 272 | } | 266 | } |
| 273 | } | 267 | } |
| 274 | | 268 | |
| 275 | /// data must outlive the returned frame | 269 | /// data must outlive the returned frame |
| 276 | pub fn preadvWindows(loop: *Loop, fd: fd_t, data: []const []u8, offset: u64) !usize { | 270 | pub fn preadvWindows(fd: fd_t, data: []const []u8, offset: u64) !usize { |
| 277 | assert(data.len != 0); | 271 | assert(data.len != 0); |
| 278 | if (data.len == 1) return preadWindows(loop, fd, data[0], offset); | 272 | if (data.len == 1) return preadWindows(fd, data[0], offset); |
| 279 | | 273 | |
| 280 | // TODO do these in parallel? | 274 | // TODO do these in parallel? |
| 281 | var off: usize = 0; | 275 | var off: usize = 0; |
| ... | @@ -283,7 +277,7 @@ pub fn preadvWindows(loop: *Loop, fd: fd_t, data: []const []u8, offset: u64) !us | ... | @@ -283,7 +277,7 @@ pub fn preadvWindows(loop: *Loop, fd: fd_t, data: []const []u8, offset: u64) !us |
| 283 | var inner_off: usize = 0; | 277 | var inner_off: usize = 0; |
| 284 | while (true) { | 278 | while (true) { |
| 285 | const v = data[iov_i]; | 279 | const v = data[iov_i]; |
| 286 | const amt_read = try preadWindows(loop, fd, v[inner_off .. v.len - inner_off], offset + off); | 280 | const amt_read = try preadWindows(fd, v[inner_off .. v.len - inner_off], offset + off); |
| 287 | off += amt_read; | 281 | off += amt_read; |
| 288 | inner_off += amt_read; | 282 | inner_off += amt_read; |
| 289 | if (inner_off == v.len) { | 283 | if (inner_off == v.len) { |
| ... | @@ -297,7 +291,7 @@ pub fn preadvWindows(loop: *Loop, fd: fd_t, data: []const []u8, offset: u64) !us | ... | @@ -297,7 +291,7 @@ pub fn preadvWindows(loop: *Loop, fd: fd_t, data: []const []u8, offset: u64) !us |
| 297 | } | 291 | } |
| 298 | } | 292 | } |
| 299 | | 293 | |
| 300 | pub fn preadWindows(loop: *Loop, fd: fd_t, data: []u8, offset: u64) !usize { | 294 | pub fn preadWindows(fd: fd_t, data: []u8, offset: u64) !usize { |
| 301 | var resume_node = Loop.ResumeNode.Basic{ | 295 | var resume_node = Loop.ResumeNode.Basic{ |
| 302 | .base = Loop.ResumeNode{ | 296 | .base = Loop.ResumeNode{ |
| 303 | .id = Loop.ResumeNode.Id.Basic, | 297 | .id = Loop.ResumeNode.Id.Basic, |
| ... | @@ -312,9 +306,9 @@ pub fn preadWindows(loop: *Loop, fd: fd_t, data: []u8, offset: u64) !usize { | ... | @@ -312,9 +306,9 @@ pub fn preadWindows(loop: *Loop, fd: fd_t, data: []u8, offset: u64) !usize { |
| 312 | }, | 306 | }, |
| 313 | }; | 307 | }; |
| 314 | // TODO only call create io completion port once per fd | 308 | // TODO only call create io completion port once per fd |
| 315 | _ = windows.CreateIoCompletionPort(fd, loop.os_data.io_port, undefined, undefined) catch undefined; | 309 | _ = windows.CreateIoCompletionPort(fd, global_event_loop.os_data.io_port, undefined, undefined) catch undefined; |
| 316 | loop.beginOneEvent(); | 310 | global_event_loop.beginOneEvent(); |
| 317 | errdefer loop.finishOneEvent(); | 311 | errdefer global_event_loop.finishOneEvent(); |
| 318 | | 312 | |
| 319 | errdefer { | 313 | errdefer { |
| 320 | _ = windows.kernel32.CancelIoEx(fd, &resume_node.base.overlapped); | 314 | _ = windows.kernel32.CancelIoEx(fd, &resume_node.base.overlapped); |
| ... | @@ -336,12 +330,7 @@ pub fn preadWindows(loop: *Loop, fd: fd_t, data: []u8, offset: u64) !usize { | ... | @@ -336,12 +330,7 @@ pub fn preadWindows(loop: *Loop, fd: fd_t, data: []u8, offset: u64) !usize { |
| 336 | } | 330 | } |
| 337 | | 331 | |
| 338 | /// iovecs must live until preadv frame completes | 332 | /// iovecs must live until preadv frame completes |
| 339 | pub fn preadvPosix( | 333 | pub fn preadvPosix(fd: fd_t, iovecs: []const os.iovec, offset: usize) os.ReadError!usize { |
| 340 | loop: *Loop, | | |
| 341 | fd: fd_t, | | |
| 342 | iovecs: []const os.iovec, | | |
| 343 | offset: usize, | | |
| 344 | ) os.ReadError!usize { | | |
| 345 | var req_node = RequestNode{ | 334 | var req_node = RequestNode{ |
| 346 | .prev = null, | 335 | .prev = null, |
| 347 | .next = null, | 336 | .next = null, |
| ... | @@ -364,21 +353,16 @@ pub fn preadvPosix( | ... | @@ -364,21 +353,16 @@ pub fn preadvPosix( |
| 364 | }, | 353 | }, |
| 365 | }; | 354 | }; |
| 366 | | 355 | |
| 367 | errdefer loop.posixFsCancel(&req_node); | 356 | errdefer global_event_loop.posixFsCancel(&req_node); |
| 368 | | 357 | |
| 369 | suspend { | 358 | suspend { |
| 370 | loop.posixFsRequest(&req_node); | 359 | global_event_loop.posixFsRequest(&req_node); |
| 371 | } | 360 | } |
| 372 | | 361 | |
| 373 | return req_node.data.msg.PReadV.result; | 362 | return req_node.data.msg.PReadV.result; |
| 374 | } | 363 | } |
| 375 | | 364 | |
| 376 | pub fn openPosix( | 365 | pub fn openPosix(path: []const u8, flags: u32, mode: File.Mode) File.OpenError!fd_t { |
| 377 | loop: *Loop, | | |
| 378 | path: []const u8, | | |
| 379 | flags: u32, | | |
| 380 | mode: File.Mode, | | |
| 381 | ) File.OpenError!fd_t { | | |
| 382 | const path_c = try std.os.toPosixPath(path); | 366 | const path_c = try std.os.toPosixPath(path); |
| 383 | | 367 | |
| 384 | var req_node = RequestNode{ | 368 | var req_node = RequestNode{ |
| ... | @@ -403,21 +387,21 @@ pub fn openPosix( | ... | @@ -403,21 +387,21 @@ pub fn openPosix( |
| 403 | }, | 387 | }, |
| 404 | }; | 388 | }; |
| 405 | | 389 | |
| 406 | errdefer loop.posixFsCancel(&req_node); | 390 | errdefer global_event_loop.posixFsCancel(&req_node); |
| 407 | | 391 | |
| 408 | suspend { | 392 | suspend { |
| 409 | loop.posixFsRequest(&req_node); | 393 | global_event_loop.posixFsRequest(&req_node); |
| 410 | } | 394 | } |
| 411 | | 395 | |
| 412 | return req_node.data.msg.Open.result; | 396 | return req_node.data.msg.Open.result; |
| 413 | } | 397 | } |
| 414 | | 398 | |
| 415 | pub fn openRead(loop: *Loop, path: []const u8) File.OpenError!fd_t { | 399 | pub fn openRead(path: []const u8) File.OpenError!fd_t { |
| 416 | switch (builtin.os) { | 400 | switch (builtin.os) { |
| 417 | .macosx, .linux, .freebsd, .netbsd, .dragonfly => { | 401 | .macosx, .linux, .freebsd, .netbsd, .dragonfly => { |
| 418 | const O_LARGEFILE = if (@hasDecl(os, "O_LARGEFILE")) os.O_LARGEFILE else 0; | 402 | const O_LARGEFILE = if (@hasDecl(os, "O_LARGEFILE")) os.O_LARGEFILE else 0; |
| 419 | const flags = O_LARGEFILE | os.O_RDONLY | os.O_CLOEXEC; | 403 | const flags = O_LARGEFILE | os.O_RDONLY | os.O_CLOEXEC; |
| 420 | return openPosix(loop, path, flags, File.default_mode); | 404 | return openPosix(path, flags, File.default_mode); |
| 421 | }, | 405 | }, |
| 422 | | 406 | |
| 423 | .windows => return windows.CreateFile( | 407 | .windows => return windows.CreateFile( |
| ... | @@ -436,12 +420,12 @@ pub fn openRead(loop: *Loop, path: []const u8) File.OpenError!fd_t { | ... | @@ -436,12 +420,12 @@ pub fn openRead(loop: *Loop, path: []const u8) File.OpenError!fd_t { |
| 436 | | 420 | |
| 437 | /// Creates if does not exist. Truncates the file if it exists. | 421 | /// Creates if does not exist. Truncates the file if it exists. |
| 438 | /// Uses the default mode. | 422 | /// Uses the default mode. |
| 439 | pub fn openWrite(loop: *Loop, path: []const u8) File.OpenError!fd_t { | 423 | pub fn openWrite(path: []const u8) File.OpenError!fd_t { |
| 440 | return openWriteMode(loop, path, File.default_mode); | 424 | return openWriteMode(path, File.default_mode); |
| 441 | } | 425 | } |
| 442 | | 426 | |
| 443 | /// Creates if does not exist. Truncates the file if it exists. | 427 | /// Creates if does not exist. Truncates the file if it exists. |
| 444 | pub fn openWriteMode(loop: *Loop, path: []const u8, mode: File.Mode) File.OpenError!fd_t { | 428 | pub fn openWriteMode(path: []const u8, mode: File.Mode) File.OpenError!fd_t { |
| 445 | switch (builtin.os) { | 429 | switch (builtin.os) { |
| 446 | .macosx, | 430 | .macosx, |
| 447 | .linux, | 431 | .linux, |
| ... | @@ -451,7 +435,7 @@ pub fn openWriteMode(loop: *Loop, path: []const u8, mode: File.Mode) File.OpenEr | ... | @@ -451,7 +435,7 @@ pub fn openWriteMode(loop: *Loop, path: []const u8, mode: File.Mode) File.OpenEr |
| 451 | => { | 435 | => { |
| 452 | const O_LARGEFILE = if (@hasDecl(os, "O_LARGEFILE")) os.O_LARGEFILE else 0; | 436 | const O_LARGEFILE = if (@hasDecl(os, "O_LARGEFILE")) os.O_LARGEFILE else 0; |
| 453 | const flags = O_LARGEFILE | os.O_WRONLY | os.O_CREAT | os.O_CLOEXEC | os.O_TRUNC; | 437 | const flags = O_LARGEFILE | os.O_WRONLY | os.O_CREAT | os.O_CLOEXEC | os.O_TRUNC; |
| 454 | return openPosix(loop, path, flags, File.default_mode); | 438 | return openPosix(path, flags, File.default_mode); |
| 455 | }, | 439 | }, |
| 456 | .windows => return windows.CreateFile( | 440 | .windows => return windows.CreateFile( |
| 457 | path, | 441 | path, |
| ... | @@ -467,16 +451,12 @@ pub fn openWriteMode(loop: *Loop, path: []const u8, mode: File.Mode) File.OpenEr | ... | @@ -467,16 +451,12 @@ pub fn openWriteMode(loop: *Loop, path: []const u8, mode: File.Mode) File.OpenEr |
| 467 | } | 451 | } |
| 468 | | 452 | |
| 469 | /// Creates if does not exist. Does not truncate. | 453 | /// Creates if does not exist. Does not truncate. |
| 470 | pub fn openReadWrite( | 454 | pub fn openReadWrite(path: []const u8, mode: File.Mode) File.OpenError!fd_t { |
| 471 | loop: *Loop, | | |
| 472 | path: []const u8, | | |
| 473 | mode: File.Mode, | | |
| 474 | ) File.OpenError!fd_t { | | |
| 475 | switch (builtin.os) { | 455 | switch (builtin.os) { |
| 476 | .macosx, .linux, .freebsd, .netbsd, .dragonfly => { | 456 | .macosx, .linux, .freebsd, .netbsd, .dragonfly => { |
| 477 | const O_LARGEFILE = if (@hasDecl(os, "O_LARGEFILE")) os.O_LARGEFILE else 0; | 457 | const O_LARGEFILE = if (@hasDecl(os, "O_LARGEFILE")) os.O_LARGEFILE else 0; |
| 478 | const flags = O_LARGEFILE | os.O_RDWR | os.O_CREAT | os.O_CLOEXEC; | 458 | const flags = O_LARGEFILE | os.O_RDWR | os.O_CREAT | os.O_CLOEXEC; |
| 479 | return openPosix(loop, path, flags, mode); | 459 | return openPosix(path, flags, mode); |
| 480 | }, | 460 | }, |
| 481 | | 461 | |
| 482 | .windows => return windows.CreateFile( | 462 | .windows => return windows.CreateFile( |
| ... | @@ -500,7 +480,7 @@ pub fn openReadWrite( | ... | @@ -500,7 +480,7 @@ pub fn openReadWrite( |
| 500 | /// If you call `setHandle` then finishing will close the fd; otherwise finishing | 480 | /// If you call `setHandle` then finishing will close the fd; otherwise finishing |
| 501 | /// will deallocate the `CloseOperation`. | 481 | /// will deallocate the `CloseOperation`. |
| 502 | pub const CloseOperation = struct { | 482 | pub const CloseOperation = struct { |
| 503 | loop: *Loop, | 483 | allocator: *Allocator, |
| 504 | os_data: OsData, | 484 | os_data: OsData, |
| 505 | | 485 | |
| 506 | const OsData = switch (builtin.os) { | 486 | const OsData = switch (builtin.os) { |
| ... | @@ -518,10 +498,10 @@ pub const CloseOperation = struct { | ... | @@ -518,10 +498,10 @@ pub const CloseOperation = struct { |
| 518 | close_req_node: RequestNode, | 498 | close_req_node: RequestNode, |
| 519 | }; | 499 | }; |
| 520 | | 500 | |
| 521 | pub fn start(loop: *Loop) (error{OutOfMemory}!*CloseOperation) { | 501 | pub fn start(allocator: *Allocator) (error{OutOfMemory}!*CloseOperation) { |
| 522 | const self = try loop.allocator.create(CloseOperation); | 502 | const self = try allocator.create(CloseOperation); |
| 523 | self.* = CloseOperation{ | 503 | self.* = CloseOperation{ |
| 524 | .loop = loop, | 504 | .allocator = allocator, |
| 525 | .os_data = switch (builtin.os) { | 505 | .os_data = switch (builtin.os) { |
| 526 | .linux, .macosx, .freebsd, .netbsd, .dragonfly => initOsDataPosix(self), | 506 | .linux, .macosx, .freebsd, .netbsd, .dragonfly => initOsDataPosix(self), |
| 527 | .windows => OsData{ .handle = null }, | 507 | .windows => OsData{ .handle = null }, |
| ... | @@ -557,16 +537,16 @@ pub const CloseOperation = struct { | ... | @@ -557,16 +537,16 @@ pub const CloseOperation = struct { |
| 557 | .dragonfly, | 537 | .dragonfly, |
| 558 | => { | 538 | => { |
| 559 | if (self.os_data.have_fd) { | 539 | if (self.os_data.have_fd) { |
| 560 | self.loop.posixFsRequest(&self.os_data.close_req_node); | 540 | global_event_loop.posixFsRequest(&self.os_data.close_req_node); |
| 561 | } else { | 541 | } else { |
| 562 | self.loop.allocator.destroy(self); | 542 | self.allocator.destroy(self); |
| 563 | } | 543 | } |
| 564 | }, | 544 | }, |
| 565 | .windows => { | 545 | .windows => { |
| 566 | if (self.os_data.handle) |handle| { | 546 | if (self.os_data.handle) |handle| { |
| 567 | os.close(handle); | 547 | os.close(handle); |
| 568 | } | 548 | } |
| 569 | self.loop.allocator.destroy(self); | 549 | self.allocator.destroy(self); |
| 570 | }, | 550 | }, |
| 571 | else => @compileError("Unsupported OS"), | 551 | else => @compileError("Unsupported OS"), |
| 572 | } | 552 | } |
| ... | @@ -629,25 +609,25 @@ pub const CloseOperation = struct { | ... | @@ -629,25 +609,25 @@ pub const CloseOperation = struct { |
| 629 | | 609 | |
| 630 | /// contents must remain alive until writeFile completes. | 610 | /// contents must remain alive until writeFile completes. |
| 631 | /// TODO make this atomic or provide writeFileAtomic and rename this one to writeFileTruncate | 611 | /// TODO make this atomic or provide writeFileAtomic and rename this one to writeFileTruncate |
| 632 | pub fn writeFile(loop: *Loop, path: []const u8, contents: []const u8) !void { | 612 | pub fn writeFile(allocator: *Allocator, path: []const u8, contents: []const u8) !void { |
| 633 | return writeFileMode(loop, path, contents, File.default_mode); | 613 | return writeFileMode(allocator, path, contents, File.default_mode); |
| 634 | } | 614 | } |
| 635 | | 615 | |
| 636 | /// contents must remain alive until writeFile completes. | 616 | /// contents must remain alive until writeFile completes. |
| 637 | pub fn writeFileMode(loop: *Loop, path: []const u8, contents: []const u8, mode: File.Mode) !void { | 617 | pub fn writeFileMode(allocator: *Allocator, path: []const u8, contents: []const u8, mode: File.Mode) !void { |
| 638 | switch (builtin.os) { | 618 | switch (builtin.os) { |
| 639 | .linux, | 619 | .linux, |
| 640 | .macosx, | 620 | .macosx, |
| 641 | .freebsd, | 621 | .freebsd, |
| 642 | .netbsd, | 622 | .netbsd, |
| 643 | .dragonfly, | 623 | .dragonfly, |
| 644 | => return writeFileModeThread(loop, path, contents, mode), | 624 | => return writeFileModeThread(allocator, path, contents, mode), |
| 645 | .windows => return writeFileWindows(loop, path, contents), | 625 | .windows => return writeFileWindows(path, contents), |
| 646 | else => @compileError("Unsupported OS"), | 626 | else => @compileError("Unsupported OS"), |
| 647 | } | 627 | } |
| 648 | } | 628 | } |
| 649 | | 629 | |
| 650 | fn writeFileWindows(loop: *Loop, path: []const u8, contents: []const u8) !void { | 630 | fn writeFileWindows(path: []const u8, contents: []const u8) !void { |
| 651 | const handle = try windows.CreateFile( | 631 | const handle = try windows.CreateFile( |
| 652 | path, | 632 | path, |
| 653 | windows.GENERIC_WRITE, | 633 | windows.GENERIC_WRITE, |
| ... | @@ -659,12 +639,12 @@ fn writeFileWindows(loop: *Loop, path: []const u8, contents: []const u8) !void { | ... | @@ -659,12 +639,12 @@ fn writeFileWindows(loop: *Loop, path: []const u8, contents: []const u8) !void { |
| 659 | ); | 639 | ); |
| 660 | defer os.close(handle); | 640 | defer os.close(handle); |
| 661 | | 641 | |
| 662 | try pwriteWindows(loop, handle, contents, 0); | 642 | try pwriteWindows(handle, contents, 0); |
| 663 | } | 643 | } |
| 664 | | 644 | |
| 665 | fn writeFileModeThread(loop: *Loop, path: []const u8, contents: []const u8, mode: File.Mode) !void { | 645 | fn writeFileModeThread(allocator: *Allocator, path: []const u8, contents: []const u8, mode: File.Mode) !void { |
| 666 | const path_with_null = try std.cstr.addNullByte(loop.allocator, path); | 646 | const path_with_null = try std.cstr.addNullByte(allocator, path); |
| 667 | defer loop.allocator.free(path_with_null); | 647 | defer allocator.free(path_with_null); |
| 668 | | 648 | |
| 669 | var req_node = RequestNode{ | 649 | var req_node = RequestNode{ |
| 670 | .prev = null, | 650 | .prev = null, |
| ... | @@ -688,10 +668,10 @@ fn writeFileModeThread(loop: *Loop, path: []const u8, contents: []const u8, mode | ... | @@ -688,10 +668,10 @@ fn writeFileModeThread(loop: *Loop, path: []const u8, contents: []const u8, mode |
| 688 | }, | 668 | }, |
| 689 | }; | 669 | }; |
| 690 | | 670 | |
| 691 | errdefer loop.posixFsCancel(&req_node); | 671 | errdefer global_event_loop.posixFsCancel(&req_node); |
| 692 | | 672 | |
| 693 | suspend { | 673 | suspend { |
| 694 | loop.posixFsRequest(&req_node); | 674 | global_event_loop.posixFsRequest(&req_node); |
| 695 | } | 675 | } |
| 696 | | 676 | |
| 697 | return req_node.data.msg.WriteFile.result; | 677 | return req_node.data.msg.WriteFile.result; |
| ... | @@ -700,21 +680,21 @@ fn writeFileModeThread(loop: *Loop, path: []const u8, contents: []const u8, mode | ... | @@ -700,21 +680,21 @@ fn writeFileModeThread(loop: *Loop, path: []const u8, contents: []const u8, mode |
| 700 | /// The frame resumes when the last data has been confirmed written, but before the file handle | 680 | /// The frame resumes when the last data has been confirmed written, but before the file handle |
| 701 | /// is closed. | 681 | /// is closed. |
| 702 | /// Caller owns returned memory. | 682 | /// Caller owns returned memory. |
| 703 | pub fn readFile(loop: *Loop, file_path: []const u8, max_size: usize) ![]u8 { | 683 | pub fn readFile(allocator: *Allocator, file_path: []const u8, max_size: usize) ![]u8 { |
| 704 | var close_op = try CloseOperation.start(loop); | 684 | var close_op = try CloseOperation.start(); |
| 705 | defer close_op.finish(); | 685 | defer close_op.finish(); |
| 706 | | 686 | |
| 707 | const fd = try openRead(loop, file_path); | 687 | const fd = try openRead(file_path); |
| 708 | close_op.setHandle(fd); | 688 | close_op.setHandle(fd); |
| 709 | | 689 | |
| 710 | var list = std.ArrayList(u8).init(loop.allocator); | 690 | var list = std.ArrayList(u8).init(allocator); |
| 711 | defer list.deinit(); | 691 | defer list.deinit(); |
| 712 | | 692 | |
| 713 | while (true) { | 693 | while (true) { |
| 714 | try list.ensureCapacity(list.len + mem.page_size); | 694 | try list.ensureCapacity(list.len + mem.page_size); |
| 715 | const buf = list.items[list.len..]; | 695 | const buf = list.items[list.len..]; |
| 716 | const buf_array = [_][]u8{buf}; | 696 | const buf_array = [_][]u8{buf}; |
| 717 | const amt = try preadv(loop, fd, buf_array, list.len); | 697 | const amt = try preadv(fd, buf_array, list.len); |
| 718 | list.len += amt; | 698 | list.len += amt; |
| 719 | if (list.len > max_size) { | 699 | if (list.len > max_size) { |
| 720 | return error.FileTooBig; | 700 | return error.FileTooBig; |
| ... | @@ -1392,16 +1372,15 @@ fn testFsWatch(loop: *Loop) !void { | ... | @@ -1392,16 +1372,15 @@ fn testFsWatch(loop: *Loop) !void { |
| 1392 | pub const OutStream = struct { | 1372 | pub const OutStream = struct { |
| 1393 | fd: fd_t, | 1373 | fd: fd_t, |
| 1394 | stream: Stream, | 1374 | stream: Stream, |
| 1395 | loop: *Loop, | 1375 | allocator: *Allocator, |
| 1396 | offset: usize, | 1376 | offset: usize, |
| 1397 | | 1377 | |
| 1398 | pub const Error = File.WriteError; | 1378 | pub const Error = File.WriteError; |
| 1399 | pub const Stream = event.io.OutStream(Error); | 1379 | pub const Stream = event.io.OutStream(Error); |
| 1400 | | 1380 | |
| 1401 | pub fn init(loop: *Loop, fd: fd_t, offset: usize) OutStream { | 1381 | pub fn init(allocator: *Allocator, fd: fd_t, offset: usize) OutStream { |
| 1402 | return OutStream{ | 1382 | return OutStream{ |
| 1403 | .fd = fd, | 1383 | .fd = fd, |
| 1404 | .loop = loop, | | |
| 1405 | .offset = offset, | 1384 | .offset = offset, |
| 1406 | .stream = Stream{ .writeFn = writeFn }, | 1385 | .stream = Stream{ .writeFn = writeFn }, |
| 1407 | }; | 1386 | }; |
| ... | @@ -1411,23 +1390,22 @@ pub const OutStream = struct { | ... | @@ -1411,23 +1390,22 @@ pub const OutStream = struct { |
| 1411 | const self = @fieldParentPtr(OutStream, "stream", out_stream); | 1390 | const self = @fieldParentPtr(OutStream, "stream", out_stream); |
| 1412 | const offset = self.offset; | 1391 | const offset = self.offset; |
| 1413 | self.offset += bytes.len; | 1392 | self.offset += bytes.len; |
| 1414 | return pwritev(self.loop, self.fd, [][]const u8{bytes}, offset); | 1393 | return pwritev(self.allocator, self.fd, [_][]const u8{bytes}, offset); |
| 1415 | } | 1394 | } |
| 1416 | }; | 1395 | }; |
| 1417 | | 1396 | |
| 1418 | pub const InStream = struct { | 1397 | pub const InStream = struct { |
| 1419 | fd: fd_t, | 1398 | fd: fd_t, |
| 1420 | stream: Stream, | 1399 | stream: Stream, |
| 1421 | loop: *Loop, | 1400 | allocator: *Allocator, |
| 1422 | offset: usize, | 1401 | offset: usize, |
| 1423 | | 1402 | |
| 1424 | pub const Error = PReadVError; // TODO make this not have OutOfMemory | 1403 | pub const Error = PReadVError; // TODO make this not have OutOfMemory |
| 1425 | pub const Stream = event.io.InStream(Error); | 1404 | pub const Stream = event.io.InStream(Error); |
| 1426 | | 1405 | |
| 1427 | pub fn init(loop: *Loop, fd: fd_t, offset: usize) InStream { | 1406 | pub fn init(allocator: *Allocator, fd: fd_t, offset: usize) InStream { |
| 1428 | return InStream{ | 1407 | return InStream{ |
| 1429 | .fd = fd, | 1408 | .fd = fd, |
| 1430 | .loop = loop, | | |
| 1431 | .offset = offset, | 1409 | .offset = offset, |
| 1432 | .stream = Stream{ .readFn = readFn }, | 1410 | .stream = Stream{ .readFn = readFn }, |
| 1433 | }; | 1411 | }; |
| ... | @@ -1435,7 +1413,7 @@ pub const InStream = struct { | ... | @@ -1435,7 +1413,7 @@ pub const InStream = struct { |
| 1435 | | 1413 | |
| 1436 | fn readFn(in_stream: *Stream, bytes: []u8) Error!usize { | 1414 | fn readFn(in_stream: *Stream, bytes: []u8) Error!usize { |
| 1437 | const self = @fieldParentPtr(InStream, "stream", in_stream); | 1415 | const self = @fieldParentPtr(InStream, "stream", in_stream); |
| 1438 | const amt = try preadv(self.loop, self.fd, [][]u8{bytes}, self.offset); | 1416 | const amt = try preadv(self.allocator, self.fd, [_][]u8{bytes}, self.offset); |
| 1439 | self.offset += amt; | 1417 | self.offset += amt; |
| 1440 | return amt; | 1418 | return amt; |
| 1441 | } | 1419 | } |