1 /++ 2 $(PITFALL 3 Please note: the api and behavior of this module is not externally stable at this time. See the documentation on specific functions for details. 4 ) 5 6 Shared core functionality including exception helpers, library loader, event loop, and possibly more. Maybe command line processor and uda helper and some basic shared annotation types. 7 8 I'll probably move the url, websocket, and ssl stuff in here too as they are often shared. Maybe a small internationalization helper type (a hook for external implementation) and COM helpers too. I might move the process helpers out to their own module - even things in here are not considered stable to library users at this time! 9 10 If you use this directly outside the arsd library despite its current instability caveats, you might consider using `static import` since names in here are likely to clash with Phobos if you use them together. `static import` will let you easily disambiguate and avoid name conflict errors if I add more here. Some names even clash deliberately to remind me to avoid some antipatterns inside the arsd modules! 11 12 ## Contributor notes 13 14 arsd.core should be focused on things that enable interoperability primarily and secondarily increased code quality between other, otherwise independent arsd modules. As a foundational library, it is not permitted to import anything outside the druntime `core` namespace, except in templates and examples not normally compiled in. This keeps it independent and avoids transitive dependency spillover to end users while also keeping compile speeds fast. To help keep builds snappy, also avoid significant use of ctfe inside this module. 15 16 On my linux computer, `dmd -unittest -main core.d` takes about a quarter second to run. We do not want this to grow. 17 18 `@safe` compatibility is ok when it isn't too big of a hassle. `@nogc` is a non-goal. I might accept it on some of the trivial functions but if it means changing the logic in any way to support, you will need a compelling argument to justify it. The arsd libs are supposed to be reliable and easy to use. That said, of course, don't be unnecessarily wasteful - if you can easily provide a reliable and easy to use way to let advanced users do their thing without hurting the other cases, let's discuss it. 19 20 If functionality is not needed by multiple existing arsd modules, consider adding a new module instead of adding it to the core. 21 22 Unittests should generally be hidden behind a special version guard so they don't interfere with end user tests. 23 24 History: 25 Added March 2023 (dub v11.0). Several functions were migrated in here at that time, noted individually. Members without a note were added with the module. 26 +/ 27 module arsd.core; 28 29 // FIXME: add callbacks on file open for tracing dependencies dynamically 30 31 // see for useful info: https://devblogs.microsoft.com/dotnet/how-async-await-really-works/ 32 33 // see: https://wiki.openssl.org/index.php/Simple_TLS_Server 34 35 ///ArsdUseCustomRuntime is used since other derived work from WebAssembly may be used and thus specified in the CLI 36 version(WebAssembly) version = ArsdUseCustomRuntime; 37 38 version(iOS) 39 { 40 version = EmptyEventLoop; 41 version = EmptyCoreEvent; 42 } 43 version(ArsdUseCustomRuntime) 44 { 45 version = EmptyEventLoop; 46 version = UseStdioWriteln; 47 } 48 else 49 { 50 version = HasFile; 51 version = HasSocket; 52 version = HasThread; 53 version = HasErrno; 54 } 55 56 version(HasThread) 57 { 58 import core.thread; 59 import core..volatile; 60 import core.atomic; 61 import core.time; 62 } 63 64 version(OSX) { 65 version(ArsdNoCocoa) 66 enum bool UseCocoa = false; 67 else 68 enum bool UseCocoa = true; 69 } 70 71 version(HasErrno) 72 import core.stdc.errno; 73 74 import core.attribute; 75 static if(!__traits(hasMember, core.attribute, "mustuse")) 76 enum mustuse; 77 78 // FIXME: add an arena allocator? can do task local destruction maybe. 79 80 // the three implementations are windows, epoll, and kqueue 81 version(Windows) { 82 version=Arsd_core_windows; 83 84 // import core.sys.windows.windows; 85 import core.sys.windows.winbase; 86 import core.sys.windows.windef; 87 import core.sys.windows.winnls; 88 import core.sys.windows.winuser; 89 import core.sys.windows.winsock2; 90 91 pragma(lib, "user32"); 92 pragma(lib, "ws2_32"); 93 } else version(linux) { 94 version=Arsd_core_epoll; 95 96 static if(__VERSION__ >= 2098) { 97 version=Arsd_core_has_cloexec; 98 } 99 } else version(FreeBSD) { 100 version=Arsd_core_kqueue; 101 102 import core.sys.freebsd.sys.event; 103 } else version(DragonFlyBSD) { 104 // NOT ACTUALLY TESTED 105 version=Arsd_core_kqueue; 106 107 import core.sys.dragonflybsd.sys.event; 108 } else version(NetBSD) { 109 // NOT ACTUALLY TESTED 110 version=Arsd_core_kqueue; 111 112 import core.sys.netbsd.sys.event; 113 } else version(OpenBSD) { 114 version=Arsd_core_kqueue; 115 116 // THIS FILE DOESN'T ACTUALLY EXIST, WE NEED TO MAKE IT 117 import core.sys.openbsd.sys.event; 118 } else version(OSX) { 119 version=Arsd_core_kqueue; 120 121 import core.sys.darwin.sys.event; 122 123 version(DigitalMars) { 124 version=OSXCocoa; 125 } 126 } 127 128 version(OSXCocoa) 129 enum CocoaAvailable = true; 130 else 131 enum CocoaAvailable = false; 132 133 version(Posix) { 134 import core.sys.posix.signal; 135 import core.sys.posix.unistd; 136 137 import core.sys.posix.sys.un; 138 import core.sys.posix.sys.socket; 139 import core.sys.posix.netinet.in_; 140 } 141 142 // FIXME: the exceptions should actually give some explanatory text too (at least sometimes) 143 144 /+ 145 ========================= 146 GENERAL UTILITY FUNCTIONS 147 ========================= 148 +/ 149 150 // enum stringz : const(char)* { init = null } 151 152 /++ 153 A wrapper around a `const(char)*` to indicate that it is a zero-terminated C string. 154 +/ 155 struct stringz { 156 private const(char)* raw; 157 158 /++ 159 Wraps the given pointer in the struct. Note that it retains a copy of the pointer. 160 +/ 161 this(const(char)* raw) { 162 this.raw = raw; 163 } 164 165 /++ 166 Returns the original raw pointer back out. 167 +/ 168 const(char)* ptr() const { 169 return raw; 170 } 171 172 /++ 173 Borrows a slice of the pointer up to (but not including) the zero terminator. 174 +/ 175 const(char)[] borrow() const { 176 if(raw is null) 177 return null; 178 179 const(char)* p = raw; 180 int length; 181 while(*p++) length++; 182 183 return raw[0 .. length]; 184 } 185 } 186 187 /++ 188 A limited variant to hold just a few types. It is made for the use of packing a small amount of extra data into error messages. 189 +/ 190 /+ 191 * if length and ptr are both 0, it is null 192 * if ptr == 1, length is an integer 193 * if ptr == 2, length is an unsigned integer (suggest printing in hex) 194 * if ptr == 3, length is a combination of flags (suggest printing in binary) 195 * if ptr == 4, length is a unix permission thing (suggest printing in octal) 196 * if ptr == 5, length is a double float 197 * if ptr == 15, length must be 0. this holds an empty, non-null, SSO string. 198 * if ptr >= 16 && < 24, length is reinterpret-casted a small string of length of (ptr & 0x7) + 1 199 * if length == size_t.max, ptr is interpreted as a stringz 200 * if ptr >= 1024, it is a non-null D string or byte array. It is a string if the length high bit is clear, a byte array if it is set. the length is what is left after you mask that out. 201 202 All other ptr values are reserved for future expansion. 203 +/ 204 struct LimitedVariant { 205 206 /++ 207 208 +/ 209 enum Contains { 210 null_, 211 intDecimal, 212 intHex, 213 intBinary, 214 intOctal, 215 double_, 216 emptySso, 217 stringSso, 218 stringz, 219 string, 220 bytes, 221 222 invalid, 223 } 224 225 /++ 226 227 +/ 228 Contains contains() const { 229 auto tag = cast(size_t) ptr; 230 if(ptr is null && length is null) 231 return Contains.null_; 232 else switch(tag) { 233 case 1: return Contains.intDecimal; 234 case 2: return Contains.intHex; 235 case 3: return Contains.intBinary; 236 case 4: return Contains.intOctal; 237 case 5: return Contains.double_; 238 case 15: return length is null ? Contains.emptySso : Contains.invalid; 239 default: 240 if(tag >= 16 && tag < 24) { 241 return Contains.stringSso; 242 } else if(tag >= 1024) { 243 if(cast(size_t) length == size_t.max) 244 return Contains.stringz; 245 else 246 return isHighBitSet ? Contains.bytes : Contains..string; 247 } else { 248 return Contains.invalid; 249 } 250 } 251 } 252 253 /// ditto 254 bool containsInt() const { 255 with(Contains) 256 switch(contains) { 257 case intDecimal, intHex, intBinary, intOctal: 258 return true; 259 default: 260 return false; 261 } 262 } 263 264 /// ditto 265 bool containsString() const { 266 with(Contains) 267 switch(contains) { 268 case null_, emptySso, stringSso, string: 269 // case stringz: 270 return true; 271 default: 272 return false; 273 } 274 } 275 276 /// ditto 277 bool containsDouble() const { 278 with(Contains) 279 switch(contains) { 280 case double_: 281 return true; 282 default: 283 return false; 284 } 285 } 286 287 /// ditto 288 bool containsBytes() const { 289 with(Contains) 290 switch(contains) { 291 case bytes, null_: 292 return true; 293 default: 294 return false; 295 } 296 } 297 298 private const(void)* length; 299 private const(ubyte)* ptr; 300 301 private void Throw() const { 302 throw ArsdException!"LimitedVariant"(cast(size_t) length, cast(size_t) ptr); 303 } 304 305 private bool isHighBitSet() const { 306 return (cast(size_t) length >> (size_t.sizeof * 8 - 1) & 0x1) != 0; 307 } 308 309 /++ 310 getString gets a reference to the string stored internally, see [toString] to get a string representation or whatever is inside. 311 312 +/ 313 const(char)[] getString() const return { 314 with(Contains) 315 switch(contains()) { 316 case null_: 317 return null; 318 case emptySso: 319 return (cast(const(char)*) ptr)[0 .. 0]; // zero length, non-null 320 case stringSso: 321 auto len = ((cast(size_t) ptr) & 0x7) + 1; 322 return (cast(char*) &length)[0 .. len]; 323 case string: 324 return (cast(const(char)*) ptr)[0 .. cast(size_t) length]; 325 default: 326 Throw(); assert(0); 327 } 328 } 329 330 /// ditto 331 long getInt() const { 332 if(containsInt) 333 return cast(long) length; 334 else 335 Throw(); 336 assert(0); 337 } 338 339 /// ditto 340 double getDouble() const { 341 if(containsDouble) 342 return *cast(double*) &length; 343 else 344 Throw(); 345 assert(0); 346 } 347 348 /// ditto 349 const(ubyte)[] getBytes() const { 350 with(Contains) 351 switch(contains()) { 352 case null_: 353 return null; 354 case bytes: 355 return ptr[0 .. (cast(size_t) length) & ((1UL << (size_t.sizeof * 8 - 1)) - 1)]; 356 default: 357 Throw(); assert(0); 358 } 359 } 360 361 /++ 362 363 +/ 364 string toString() const { 365 366 string intHelper(string prefix, int radix) { 367 char[128] buffer; 368 buffer[0 .. prefix.length] = prefix[]; 369 char[] toUse = buffer[prefix.length .. $]; 370 371 auto got = intToString(getInt(), toUse[], IntToStringArgs().withRadix(radix)); 372 373 return buffer[0 .. prefix.length + got.length].idup; 374 } 375 376 with(Contains) 377 final switch(contains()) { 378 case null_: 379 return "<null>"; 380 case intDecimal: 381 return intHelper("", 10); 382 case intHex: 383 return intHelper("0x", 16); 384 case intBinary: 385 return intHelper("0b", 2); 386 case intOctal: 387 return intHelper("0o", 8); 388 case emptySso, stringSso, string: 389 return getString().idup; 390 case bytes: 391 auto b = getBytes(); 392 393 return "<bytes>"; // FIXME 394 395 case double_: 396 assert(0); // FIXME 397 case stringz: 398 assert(0); // FIXME 399 case invalid: 400 return "<invalid>"; 401 } 402 } 403 404 /++ 405 406 +/ 407 this(string s) { 408 ptr = cast(const(ubyte)*) s.ptr; 409 length = cast(void*) s.length; 410 } 411 412 /// ditto 413 this(const(ubyte)[] b) { 414 ptr = cast(const(ubyte)*) b.ptr; 415 length = cast(void*) (b.length | (1UL << (size_t.sizeof * 8 - 1))); 416 } 417 418 /// ditto 419 this(long l, int base = 10) { 420 int tag; 421 switch(base) { 422 case 10: tag = 1; break; 423 case 16: tag = 2; break; 424 case 2: tag = 3; break; 425 case 8: tag = 4; break; 426 default: assert(0, "You passed an invalid base to LimitedVariant"); 427 } 428 ptr = cast(ubyte*) tag; 429 length = cast(void*) l; 430 } 431 432 /// ditto 433 version(none) 434 this(double d) { 435 // this crashes dmd! omg 436 assert(0); 437 // ptr = cast(ubyte*) 15; 438 // length = cast(void*) *cast(size_t*) &d; 439 } 440 } 441 442 unittest { 443 LimitedVariant v = LimitedVariant("foo"); 444 assert(v.containsString()); 445 assert(!v.containsInt()); 446 assert(v.getString() == "foo"); 447 448 LimitedVariant v2 = LimitedVariant(4); 449 assert(v2.containsInt()); 450 assert(!v2.containsString()); 451 assert(v2.getInt() == 4); 452 453 LimitedVariant v3 = LimitedVariant(cast(ubyte[]) [1, 2, 3]); 454 assert(v3.containsBytes()); 455 assert(!v3.containsString()); 456 assert(v3.getBytes() == [1, 2, 3]); 457 } 458 459 /++ 460 This is a dummy type to indicate the end of normal arguments and the beginning of the file/line inferred args. It is meant to ensure you don't accidentally send a string that is interpreted as a filename when it was meant to be a normal argument to the function and trigger the wrong overload. 461 +/ 462 struct ArgSentinel {} 463 464 /++ 465 A trivial wrapper around C's malloc that creates a D slice. It multiples n by T.sizeof and returns the slice of the pointer from 0 to n. 466 467 Please note that the ptr might be null - it is your responsibility to check that, same as normal malloc. Check `ret is null` specifically, since `ret.length` will always be `n`, even if the `malloc` failed. 468 469 Remember to `free` the returned pointer with `core.stdc.stdlib.free(ret.ptr);` 470 471 $(TIP 472 I strongly recommend you simply use the normal garbage collector unless you have a very specific reason not to. 473 ) 474 475 See_Also: 476 [mallocedStringz] 477 +/ 478 T[] mallocSlice(T)(size_t n) { 479 import c = core.stdc.stdlib; 480 481 return (cast(T*) c.malloc(n * T.sizeof))[0 .. n]; 482 } 483 484 /++ 485 Uses C's malloc to allocate a copy of `original` with an attached zero terminator. It may return a slice with a `null` pointer (but non-zero length!) if `malloc` fails and you are responsible for freeing the returned pointer with `core.stdc.stdlib.free(ret.ptr)`. 486 487 $(TIP 488 I strongly recommend you use [CharzBuffer] or Phobos' [std.string.toStringz] instead unless there's a special reason not to. 489 ) 490 491 See_Also: 492 [CharzBuffer] for a generally better alternative. You should only use `mallocedStringz` where `CharzBuffer` cannot be used (e.g. when druntime is not usable or you have no stack space for the temporary buffer). 493 494 [mallocSlice] is the function this function calls, so the notes in its documentation applies here too. 495 +/ 496 char[] mallocedStringz(in char[] original) { 497 auto slice = mallocSlice!char(original.length + 1); 498 if(slice is null) 499 return null; 500 slice[0 .. original.length] = original[]; 501 slice[original.length] = 0; 502 return slice; 503 } 504 505 /++ 506 Basically a `scope class` you can return from a function or embed in another aggregate. 507 +/ 508 struct OwnedClass(Class) { 509 ubyte[__traits(classInstanceSize, Class)] rawData; 510 511 static OwnedClass!Class defaultConstructed() { 512 OwnedClass!Class i = OwnedClass!Class.init; 513 i.initializeRawData(); 514 return i; 515 } 516 517 private void initializeRawData() @trusted { 518 if(!this) 519 rawData[] = cast(ubyte[]) typeid(Class).initializer[]; 520 } 521 522 this(T...)(T t) { 523 initializeRawData(); 524 rawInstance.__ctor(t); 525 } 526 527 bool opCast(T : bool)() @trusted { 528 return !(*(cast(void**) rawData.ptr) is null); 529 } 530 531 @disable this(); 532 @disable this(this); 533 534 Class rawInstance() return @trusted { 535 if(!this) 536 throw new Exception("null"); 537 return cast(Class) rawData.ptr; 538 } 539 540 alias rawInstance this; 541 542 ~this() @trusted { 543 if(this) 544 .destroy(rawInstance()); 545 } 546 } 547 548 // might move RecyclableMemory here 549 550 version(Posix) 551 package(arsd) void makeNonBlocking(int fd) { 552 import core.sys.posix.fcntl; 553 auto flags = fcntl(fd, F_GETFL, 0); 554 if(flags == -1) 555 throw new ErrnoApiException("fcntl get", errno); 556 flags |= O_NONBLOCK; 557 auto s = fcntl(fd, F_SETFL, flags); 558 if(s == -1) 559 throw new ErrnoApiException("fcntl set", errno); 560 } 561 562 version(Posix) 563 package(arsd) void setCloExec(int fd) { 564 import core.sys.posix.fcntl; 565 auto flags = fcntl(fd, F_GETFD, 0); 566 if(flags == -1) 567 throw new ErrnoApiException("fcntl get", errno); 568 flags |= FD_CLOEXEC; 569 auto s = fcntl(fd, F_SETFD, flags); 570 if(s == -1) 571 throw new ErrnoApiException("fcntl set", errno); 572 } 573 574 575 /++ 576 A helper object for temporarily constructing a string appropriate for the Windows API from a D UTF-8 string. 577 578 579 It will use a small internal static buffer is possible, and allocate a new buffer if the string is too big. 580 581 History: 582 Moved from simpledisplay.d to core.d in March 2023 (dub v11.0). 583 +/ 584 version(Windows) 585 struct WCharzBuffer { 586 private wchar[] buffer; 587 private wchar[128] staticBuffer = void; 588 589 /// Length of the string, excluding the zero terminator. 590 size_t length() { 591 return buffer.length; 592 } 593 594 // Returns the pointer to the internal buffer. You must assume its lifetime is less than that of the WCharzBuffer. It is zero-terminated. 595 wchar* ptr() { 596 return buffer.ptr; 597 } 598 599 /// Returns the slice of the internal buffer, excluding the zero terminator (though there is one present right off the end of the slice). You must assume its lifetime is less than that of the WCharzBuffer. 600 wchar[] slice() { 601 return buffer; 602 } 603 604 /// Copies it into a static array of wchars 605 void copyInto(R)(ref R r) { 606 static if(is(R == wchar[N], size_t N)) { 607 r[0 .. this.length] = slice[]; 608 r[this.length] = 0; 609 } else static assert(0, "can only copy into wchar[n], not " ~ R.stringof); 610 } 611 612 /++ 613 conversionFlags = [WindowsStringConversionFlags] 614 +/ 615 this(in char[] data, int conversionFlags = 0) { 616 conversionFlags |= WindowsStringConversionFlags.zeroTerminate; // this ALWAYS zero terminates cuz of its name 617 auto sz = sizeOfConvertedWstring(data, conversionFlags); 618 if(sz > staticBuffer.length) 619 buffer = new wchar[](sz); 620 else 621 buffer = staticBuffer[]; 622 623 buffer = makeWindowsString(data, buffer, conversionFlags); 624 } 625 } 626 627 /++ 628 Alternative for toStringz 629 630 History: 631 Added March 18, 2023 (dub v11.0) 632 +/ 633 struct CharzBuffer { 634 private char[] buffer; 635 private char[128] staticBuffer = void; 636 637 /// Length of the string, excluding the zero terminator. 638 size_t length() { 639 assert(buffer.length > 0); 640 return buffer.length - 1; 641 } 642 643 // Returns the pointer to the internal buffer. You must assume its lifetime is less than that of the CharzBuffer. It is zero-terminated. 644 char* ptr() { 645 return buffer.ptr; 646 } 647 648 /// Returns the slice of the internal buffer, excluding the zero terminator (though there is one present right off the end of the slice). You must assume its lifetime is less than that of the CharzBuffer. 649 char[] slice() { 650 assert(buffer.length > 0); 651 return buffer[0 .. $-1]; 652 } 653 654 /// Copies it into a static array of chars 655 void copyInto(R)(ref R r) { 656 static if(is(R == char[N], size_t N)) { 657 r[0 .. this.length] = slice[]; 658 r[this.length] = 0; 659 } else static assert(0, "can only copy into char[n], not " ~ R.stringof); 660 } 661 662 @disable this(); 663 @disable this(this); 664 665 /++ 666 Copies `data` into the CharzBuffer, allocating a new one if needed, and zero-terminates it. 667 +/ 668 this(in char[] data) { 669 if(data.length + 1 > staticBuffer.length) 670 buffer = new char[](data.length + 1); 671 else 672 buffer = staticBuffer[]; 673 674 buffer[0 .. data.length] = data[]; 675 buffer[data.length] = 0; 676 } 677 } 678 679 /++ 680 Given the string `str`, converts it to a string compatible with the Windows API and puts the result in `buffer`, returning the slice of `buffer` actually used. `buffer` must be at least [sizeOfConvertedWstring] elements long. 681 682 History: 683 Moved from simpledisplay.d to core.d in March 2023 (dub v11.0). 684 +/ 685 version(Windows) 686 wchar[] makeWindowsString(in char[] str, wchar[] buffer, int conversionFlags = WindowsStringConversionFlags.zeroTerminate) { 687 if(str.length == 0) 688 return null; 689 690 int pos = 0; 691 dchar last; 692 foreach(dchar c; str) { 693 if(c <= 0xFFFF) { 694 if((conversionFlags & WindowsStringConversionFlags.convertNewLines) && c == 10 && last != 13) 695 buffer[pos++] = 13; 696 buffer[pos++] = cast(wchar) c; 697 } else if(c <= 0x10FFFF) { 698 buffer[pos++] = cast(wchar)((((c - 0x10000) >> 10) & 0x3FF) + 0xD800); 699 buffer[pos++] = cast(wchar)(((c - 0x10000) & 0x3FF) + 0xDC00); 700 } 701 702 last = c; 703 } 704 705 if(conversionFlags & WindowsStringConversionFlags.zeroTerminate) { 706 buffer[pos] = 0; 707 } 708 709 return buffer[0 .. pos]; 710 } 711 712 /++ 713 Converts the Windows API string `str` to a D UTF-8 string, storing it in `buffer`. Returns the slice of `buffer` actually used. 714 715 History: 716 Moved from simpledisplay.d to core.d in March 2023 (dub v11.0). 717 +/ 718 version(Windows) 719 char[] makeUtf8StringFromWindowsString(in wchar[] str, char[] buffer) { 720 if(str.length == 0) 721 return null; 722 723 auto got = WideCharToMultiByte(CP_UTF8, 0, str.ptr, cast(int) str.length, buffer.ptr, cast(int) buffer.length, null, null); 724 if(got == 0) { 725 if(GetLastError() == ERROR_INSUFFICIENT_BUFFER) 726 throw new object.Exception("not enough buffer"); 727 else 728 throw new object.Exception("conversion"); // FIXME: GetLastError 729 } 730 return buffer[0 .. got]; 731 } 732 733 /++ 734 Converts the Windows API string `str` to a newly-allocated D UTF-8 string. 735 736 History: 737 Moved from simpledisplay.d to core.d in March 2023 (dub v11.0). 738 +/ 739 version(Windows) 740 string makeUtf8StringFromWindowsString(in wchar[] str) { 741 char[] buffer; 742 auto got = WideCharToMultiByte(CP_UTF8, 0, str.ptr, cast(int) str.length, null, 0, null, null); 743 buffer.length = got; 744 745 // it is unique because we just allocated it above! 746 return cast(string) makeUtf8StringFromWindowsString(str, buffer); 747 } 748 749 /// ditto 750 version(Windows) 751 string makeUtf8StringFromWindowsString(wchar* str) { 752 char[] buffer; 753 auto got = WideCharToMultiByte(CP_UTF8, 0, str, -1, null, 0, null, null); 754 buffer.length = got; 755 756 got = WideCharToMultiByte(CP_UTF8, 0, str, -1, buffer.ptr, cast(int) buffer.length, null, null); 757 if(got == 0) { 758 if(GetLastError() == ERROR_INSUFFICIENT_BUFFER) 759 throw new object.Exception("not enough buffer"); 760 else 761 throw new object.Exception("conversion"); // FIXME: GetLastError 762 } 763 return cast(string) buffer[0 .. got]; 764 } 765 766 // only used from minigui rn 767 package int findIndexOfZero(in wchar[] str) { 768 foreach(idx, wchar ch; str) 769 if(ch == 0) 770 return cast(int) idx; 771 return cast(int) str.length; 772 } 773 package int findIndexOfZero(in char[] str) { 774 foreach(idx, char ch; str) 775 if(ch == 0) 776 return cast(int) idx; 777 return cast(int) str.length; 778 } 779 780 /++ 781 Returns a minimum buffer length to hold the string `s` with the given conversions. It might be slightly larger than necessary, but is guaranteed to be big enough to hold it. 782 783 History: 784 Moved from simpledisplay.d to core.d in March 2023 (dub v11.0). 785 +/ 786 version(Windows) 787 int sizeOfConvertedWstring(in char[] s, int conversionFlags) { 788 int size = 0; 789 790 if(conversionFlags & WindowsStringConversionFlags.convertNewLines) { 791 // need to convert line endings, which means the length will get bigger. 792 793 // BTW I betcha this could be faster with some simd stuff. 794 char last; 795 foreach(char ch; s) { 796 if(ch == 10 && last != 13) 797 size++; // will add a 13 before it... 798 size++; 799 last = ch; 800 } 801 } else { 802 // no conversion necessary, just estimate based on length 803 /* 804 I don't think there's any string with a longer length 805 in code units when encoded in UTF-16 than it has in UTF-8. 806 This will probably over allocate, but that's OK. 807 */ 808 size = cast(int) s.length; 809 } 810 811 if(conversionFlags & WindowsStringConversionFlags.zeroTerminate) 812 size++; 813 814 return size; 815 } 816 817 /++ 818 Used by [makeWindowsString] and [WCharzBuffer] 819 820 History: 821 Moved from simpledisplay.d to core.d in March 2023 (dub v11.0). 822 +/ 823 version(Windows) 824 enum WindowsStringConversionFlags : int { 825 /++ 826 Append a zero terminator to the string. 827 +/ 828 zeroTerminate = 1, 829 /++ 830 Converts newlines from \n to \r\n. 831 +/ 832 convertNewLines = 2, 833 } 834 835 /++ 836 An int printing function that doesn't need to import Phobos. Can do some of the things std.conv.to and std.format.format do. 837 838 The buffer must be sized to hold the converted number. 32 chars is enough for most anything. 839 840 Returns: the slice of `buffer` containing the converted number. 841 +/ 842 char[] intToString(long value, char[] buffer, IntToStringArgs args = IntToStringArgs.init) { 843 const int radix = args.radix ? args.radix : 10; 844 const int digitsPad = args.padTo; 845 const int groupSize = args.groupSize; 846 847 int pos; 848 849 if(value < 0) { 850 buffer[pos++] = '-'; 851 value = -value; 852 } 853 854 int start = pos; 855 int digitCount; 856 857 do { 858 auto remainder = value % radix; 859 value = value / radix; 860 861 buffer[pos++] = cast(char) (remainder < 10 ? (remainder + '0') : (remainder - 10 + args.ten)); 862 digitCount++; 863 } while(value); 864 865 if(digitsPad > 0) { 866 while(digitCount < digitsPad) { 867 buffer[pos++] = args.padWith; 868 digitCount++; 869 } 870 } 871 872 assert(pos >= 1); 873 assert(pos - start > 0); 874 875 auto reverseSlice = buffer[start .. pos]; 876 for(int i = 0; i < reverseSlice.length / 2; i++) { 877 auto paired = cast(int) reverseSlice.length - i - 1; 878 char tmp = reverseSlice[i]; 879 reverseSlice[i] = reverseSlice[paired]; 880 reverseSlice[paired] = tmp; 881 } 882 883 return buffer[0 .. pos]; 884 } 885 886 /// ditto 887 struct IntToStringArgs { 888 private { 889 ubyte padTo; 890 char padWith; 891 ubyte radix; 892 char ten; 893 ubyte groupSize; 894 char separator; 895 } 896 897 IntToStringArgs withPadding(int padTo, char padWith = '0') { 898 IntToStringArgs args = this; 899 args.padTo = cast(ubyte) padTo; 900 args.padWith = padWith; 901 return args; 902 } 903 904 IntToStringArgs withRadix(int radix, char ten = 'a') { 905 IntToStringArgs args = this; 906 args.radix = cast(ubyte) radix; 907 args.ten = ten; 908 return args; 909 } 910 911 IntToStringArgs withGroupSeparator(int groupSize, char separator = '_') { 912 IntToStringArgs args = this; 913 args.groupSize = cast(ubyte) groupSize; 914 args.separator = separator; 915 return args; 916 } 917 } 918 919 unittest { 920 char[32] buffer; 921 assert(intToString(0, buffer[]) == "0"); 922 assert(intToString(-1, buffer[]) == "-1"); 923 assert(intToString(-132, buffer[]) == "-132"); 924 assert(intToString(-1932, buffer[]) == "-1932"); 925 assert(intToString(1, buffer[]) == "1"); 926 assert(intToString(132, buffer[]) == "132"); 927 assert(intToString(1932, buffer[]) == "1932"); 928 929 assert(intToString(0x1, buffer[], IntToStringArgs().withRadix(16)) == "1"); 930 assert(intToString(0x1b, buffer[], IntToStringArgs().withRadix(16)) == "1b"); 931 assert(intToString(0xef1, buffer[], IntToStringArgs().withRadix(16)) == "ef1"); 932 933 assert(intToString(0xef1, buffer[], IntToStringArgs().withRadix(16).withPadding(8)) == "00000ef1"); 934 assert(intToString(-0xef1, buffer[], IntToStringArgs().withRadix(16).withPadding(8)) == "-00000ef1"); 935 assert(intToString(-0xef1, buffer[], IntToStringArgs().withRadix(16, 'A').withPadding(8, ' ')) == "- EF1"); 936 } 937 938 /++ 939 History: 940 Moved from color.d to core.d in March 2023 (dub v11.0). 941 +/ 942 nothrow @safe @nogc pure 943 inout(char)[] stripInternal(return inout(char)[] s) { 944 bool isAllWhitespace = true; 945 foreach(i, char c; s) 946 if(c != ' ' && c != '\t' && c != '\n' && c != '\r') { 947 s = s[i .. $]; 948 isAllWhitespace = false; 949 break; 950 } 951 952 if(isAllWhitespace) 953 return s[$..$]; 954 955 for(int a = cast(int)(s.length - 1); a > 0; a--) { 956 char c = s[a]; 957 if(c != ' ' && c != '\t' && c != '\n' && c != '\r') { 958 s = s[0 .. a + 1]; 959 break; 960 } 961 } 962 963 return s; 964 } 965 966 /// ditto 967 nothrow @safe @nogc pure 968 inout(char)[] stripRightInternal(return inout(char)[] s) { 969 bool isAllWhitespace = true; 970 foreach_reverse(a, c; s) { 971 if(c != ' ' && c != '\t' && c != '\n' && c != '\r') { 972 s = s[0 .. a + 1]; 973 isAllWhitespace = false; 974 break; 975 } 976 } 977 if(isAllWhitespace) 978 s = s[0..0]; 979 980 return s; 981 982 } 983 984 /++ 985 Shortcut for converting some types to string without invoking Phobos (but it will as a last resort). 986 987 History: 988 Moved from color.d to core.d in March 2023 (dub v11.0). 989 +/ 990 string toStringInternal(T)(T t) { 991 char[32] buffer; 992 static if(is(T : string)) 993 return t; 994 else static if(is(T : long)) 995 return intToString(t, buffer[]).idup; 996 else static if(is(T == enum)) { 997 switch(t) { 998 foreach(memberName; __traits(allMembers, T)) { 999 case __traits(getMember, T, memberName): 1000 return memberName; 1001 } 1002 default: 1003 return "<unknown>"; 1004 } 1005 } else { 1006 import std.conv; 1007 return to!string(t); 1008 } 1009 } 1010 1011 /++ 1012 1013 +/ 1014 string flagsToString(Flags)(ulong value) { 1015 string r; 1016 1017 void add(string memberName) { 1018 if(r.length) 1019 r ~= " | "; 1020 r ~= memberName; 1021 } 1022 1023 string none = "<none>"; 1024 1025 foreach(memberName; __traits(allMembers, Flags)) { 1026 auto flag = cast(ulong) __traits(getMember, Flags, memberName); 1027 if(flag) { 1028 if((value & flag) == flag) 1029 add(memberName); 1030 } else { 1031 none = memberName; 1032 } 1033 } 1034 1035 if(r.length == 0) 1036 r = none; 1037 1038 return r; 1039 } 1040 1041 unittest { 1042 enum MyFlags { 1043 none = 0, 1044 a = 1, 1045 b = 2 1046 } 1047 1048 assert(flagsToString!MyFlags(3) == "a | b"); 1049 assert(flagsToString!MyFlags(0) == "none"); 1050 assert(flagsToString!MyFlags(2) == "b"); 1051 } 1052 1053 /++ 1054 This populates a struct from a list of values (or other expressions, but it only looks at the values) based on types of the members, with one exception: `bool` members.. maybe. 1055 1056 It is intended for collecting a record of relevant UDAs off a symbol in a single call like this: 1057 1058 --- 1059 struct Name { 1060 string n; 1061 } 1062 1063 struct Validator { 1064 string regex; 1065 } 1066 1067 struct FormInfo { 1068 Name name; 1069 Validator validator; 1070 } 1071 1072 @Name("foo") @Validator(".*") 1073 void foo() {} 1074 1075 auto info = populateFromUdas!(FormInfo, __traits(getAttributes, foo)); 1076 assert(info.name == Name("foo")); 1077 assert(info.validator == Validator(".*")); 1078 --- 1079 1080 Note that instead of UDAs, you can also pass a variadic argument list and get the same result, but the function is `populateFromArgs` and you pass them as the runtime list to bypass "args cannot be evaluated at compile time" errors: 1081 1082 --- 1083 void foo(T...)(T t) { 1084 auto info = populateFromArgs!(FormInfo)(t); 1085 // assuming the call below 1086 assert(info.name == Name("foo")); 1087 assert(info.validator == Validator(".*")); 1088 } 1089 1090 foo(Name("foo"), Validator(".*")); 1091 --- 1092 1093 The benefit of this over constructing the struct directly is that the arguments can be reordered or missing. Its value is diminished with named arguments in the language. 1094 +/ 1095 template populateFromUdas(Struct, UDAs...) { 1096 enum Struct populateFromUdas = () { 1097 Struct ret; 1098 foreach(memberName; __traits(allMembers, Struct)) { 1099 alias memberType = typeof(__traits(getMember, Struct, memberName)); 1100 foreach(uda; UDAs) { 1101 static if(is(memberType == PresenceOf!a, a)) { 1102 static if(__traits(isSame, a, uda)) 1103 __traits(getMember, ret, memberName) = true; 1104 } 1105 else 1106 static if(is(typeof(uda) : memberType)) { 1107 __traits(getMember, ret, memberName) = uda; 1108 } 1109 } 1110 } 1111 1112 return ret; 1113 }(); 1114 } 1115 1116 /// ditto 1117 Struct populateFromArgs(Struct, Args...)(Args args) { 1118 Struct ret; 1119 foreach(memberName; __traits(allMembers, Struct)) { 1120 alias memberType = typeof(__traits(getMember, Struct, memberName)); 1121 foreach(arg; args) { 1122 static if(is(typeof(arg == memberType))) { 1123 __traits(getMember, ret, memberName) = arg; 1124 } 1125 } 1126 } 1127 1128 return ret; 1129 } 1130 1131 /// ditto 1132 struct PresenceOf(alias a) { 1133 bool there; 1134 alias there this; 1135 } 1136 1137 /// 1138 unittest { 1139 enum a; 1140 enum b; 1141 struct Name { string name; } 1142 struct Info { 1143 Name n; 1144 PresenceOf!a athere; 1145 PresenceOf!b bthere; 1146 int c; 1147 } 1148 1149 void test() @a @Name("test") {} 1150 1151 auto info = populateFromUdas!(Info, __traits(getAttributes, test)); 1152 assert(info.n == Name("test")); // but present ones are in there 1153 assert(info.athere == true); // non-values can be tested with PresenceOf!it, which works like a bool 1154 assert(info.bthere == false); 1155 assert(info.c == 0); // absent thing will keep the default value 1156 } 1157 1158 /++ 1159 Declares a delegate property with several setters to allow for handlers that don't care about the arguments. 1160 1161 Throughout the arsd library, you will often see types of these to indicate that you can set listeners with or without arguments. If you care about the details of the callback event, you can set a delegate that declares them. And if you don't, you can set one that doesn't even declare them and it will be ignored. 1162 +/ 1163 struct FlexibleDelegate(DelegateType) { 1164 // please note that Parameters and ReturnType are public now! 1165 static if(is(DelegateType FunctionType == delegate)) 1166 static if(is(FunctionType Parameters == __parameters)) 1167 static if(is(DelegateType ReturnType == return)) { 1168 1169 /++ 1170 Calls the currently set delegate. 1171 1172 Diagnostics: 1173 If the callback delegate has not been set, this may cause a null pointer dereference. 1174 +/ 1175 ReturnType opCall(Parameters args) { 1176 return dg(args); 1177 } 1178 1179 /++ 1180 Use `if(thing)` to check if the delegate is null or not. 1181 +/ 1182 bool opCast(T : bool)() { 1183 return dg !is null; 1184 } 1185 1186 /++ 1187 These opAssign overloads are what puts the flexibility in the flexible delegate. 1188 1189 Bugs: 1190 The other overloads do not keep attributes like `nothrow` on the `dg` parameter, making them unusable if `DelegateType` requires them. I consider the attributes more trouble than they're worth anyway, and the language's poor support for composing them doesn't help any. I have no need for them and thus no plans to add them in the overloads at this time. 1191 +/ 1192 void opAssign(DelegateType dg) { 1193 this.dg = dg; 1194 } 1195 1196 /// ditto 1197 void opAssign(ReturnType delegate() dg) { 1198 this.dg = (Parameters ignored) => dg(); 1199 } 1200 1201 /// ditto 1202 void opAssign(ReturnType function(Parameters params) dg) { 1203 this.dg = (Parameters params) => dg(params); 1204 } 1205 1206 /// ditto 1207 void opAssign(ReturnType function() dg) { 1208 this.dg = (Parameters ignored) => dg(); 1209 } 1210 1211 /// ditto 1212 void opAssign(typeof(null) explicitNull) { 1213 this.dg = null; 1214 } 1215 1216 private DelegateType dg; 1217 } 1218 else static assert(0, DelegateType.stringof ~ " failed return value check"); 1219 else static assert(0, DelegateType.stringof ~ " failed parameters check"); 1220 else static assert(0, DelegateType.stringof ~ " failed delegate check"); 1221 } 1222 1223 /++ 1224 1225 +/ 1226 unittest { 1227 // you don't have to put the arguments in a struct, but i recommend 1228 // you do as it is more future proof - you can add more info to the 1229 // struct without breaking user code that consumes it. 1230 struct MyEventArguments { 1231 1232 } 1233 1234 // then you declare it just adding FlexibleDelegate!() around the 1235 // plain delegate type you'd normally use 1236 FlexibleDelegate!(void delegate(MyEventArguments args)) callback; 1237 1238 // until you set it, it will be null and thus be false in any boolean check 1239 assert(!callback); 1240 1241 // can set it to the properly typed thing 1242 callback = delegate(MyEventArguments args) {}; 1243 1244 // and now it is no longer null 1245 assert(callback); 1246 1247 // or if you don't care about the args, you can leave them off 1248 callback = () {}; 1249 1250 // and it works if the compiler types you as a function instead of delegate too 1251 // (which happens automatically if you don't access any local state or if you 1252 // explicitly define it as a function) 1253 1254 callback = function(MyEventArguments args) { }; 1255 1256 // can set it back to null explicitly if you ever wanted 1257 callback = null; 1258 1259 // the reflection info used internally also happens to be exposed publicly 1260 // which can actually sometimes be nice so if the language changes, i'll change 1261 // the code to keep this working. 1262 static assert(is(callback.ReturnType == void)); 1263 1264 // which can be convenient if the params is an annoying type since you can 1265 // consistently use something like this too 1266 callback = (callback.Parameters params) {}; 1267 1268 // check for null and call it pretty normally 1269 if(callback) 1270 callback(MyEventArguments()); 1271 } 1272 1273 /+ 1274 ====================== 1275 ERROR HANDLING HELPERS 1276 ====================== 1277 +/ 1278 1279 /+ + 1280 arsd code shouldn't be using Exception. Really, I don't think any code should be - instead, construct an appropriate object with structured information. 1281 1282 If you want to catch someone else's Exception, use `catch(object.Exception e)`. 1283 +/ 1284 //package deprecated struct Exception {} 1285 1286 1287 /++ 1288 Base class representing my exceptions. You should almost never work with this directly, but you might catch it as a generic thing. Catch it before generic `object.Exception` or `object.Throwable` in any catch chains. 1289 1290 1291 $(H3 General guidelines for exceptions) 1292 1293 The purpose of an exception is to cancel a task that has proven to be impossible and give the programmer enough information to use at a higher level to decide what to do about it. 1294 1295 Cancelling a task is accomplished with the `throw` keyword. The transmission of information to a higher level is done by the language runtime. The decision point is marked by the `catch` keyword. The part missing - the job of the `Exception` class you construct and throw - is to gather the information that will be useful at a later decision point. 1296 1297 It is thus important that you gather as much useful information as possible and keep it in a way that the code catching the exception can still interpret it when constructing an exception. Other concerns are secondary to this to this primary goal. 1298 1299 With this in mind, here's some guidelines for exception handling in arsd code. 1300 1301 $(H4 Allocations and lifetimes) 1302 1303 Don't get clever with exception allocations. You don't know what the catcher is going to do with an exception and you don't want the error handling scheme to introduce its own tricky bugs. Remember, an exception object's first job is to deliver useful information up the call chain in a way this code can use it. You don't know what this code is or what it is going to do. 1304 1305 Keep your memory management schemes simple and let the garbage collector do its job. 1306 1307 $(LIST 1308 * All thrown exceptions should be allocated with the `new` keyword. 1309 1310 * Members inside the exception should be value types or have infinite lifetime (that is, be GC managed). 1311 1312 * While this document is concerned with throwing, you might want to add additional information to an in-flight exception, and this is done by catching, so you need to know how that works too, and there is a global compiler switch that can change things, so even inside arsd we can't completely avoid its implications. 1313 1314 DIP1008's presence complicates things a bit on the catch side - if you catch an exception and return it from a function, remember to `ex.refcount = ex.refcount + 1;` so you don't introduce more use-after-free woes for those unfortunate souls. 1315 ) 1316 1317 $(H4 Error strings) 1318 1319 Strings can deliver useful information to people reading the message, but are often suboptimal for delivering useful information to other chunks of code. Remember, an exception's first job is to be caught by another block of code. Printing to users is a last resort; even if you want a user-readable error message, an exception is not the ideal way to deliver one since it is constructed in the guts of a failed task, without the higher level context of what the user was actually trying to do. User error messages ought to be made from information in the exception, combined with higher level knowledge. This is best done in a `catch` block, not a `throw` statement. 1320 1321 As such, I recommend that you: 1322 1323 $(LIST 1324 * Don't concatenate error strings at the throw site. Instead, pass the data you would have used to build the string as actual data to the constructor. This lets catchers see the original data without having to try to extract it from a string. For unique data, you will likely need a unique exception type. More on this in the next section. 1325 1326 * Don't construct error strings in a constructor either, for the same reason. Pass the useful data up the call chain, as exception members, to the maximum extent possible. Exception: if you are passed some data with a temporary lifetime that is important enough to pass up the chain. You may `.idup` or `to!string` to preserve as much data as you can before it is lost, but still store it in a separate member of the Exception subclass object. 1327 1328 * $(I Do) construct strings out of public members in [getAdditionalPrintableInformation]. When this is called, the user has requested as much relevant information as reasonable in string format. Still, avoid concatenation - it lets you pass as many key/value pairs as you like to the caller. They can concatenate as needed. However, note the words "public members" - everything you do in `getAdditionalPrintableInformation` ought to also be possible for code that caught your exception via your public methods and properties. 1329 ) 1330 1331 $(H4 Subclasses) 1332 1333 Any exception with unique data types should be a unique class. Whenever practical, this should be one you write and document at the top-level of a module. But I know we get lazy - me too - and this is why in standard D we'd often fall back to `throw new Exception("some string " ~ some info)`. To help resist these urges, I offer some helper functions to use instead that better achieve the key goal of exceptions - passing structured data up a call chain - while still being convenient to write. 1334 1335 See: [ArsdException], [Win32Enforce] 1336 1337 +/ 1338 class ArsdExceptionBase : object.Exception { 1339 /++ 1340 Don't call this except from other exceptions; this is essentially an abstract class. 1341 1342 Params: 1343 operation = the specific operation that failed, throwing the exception 1344 +/ 1345 package this(string operation, string file = __FILE__, size_t line = __LINE__, Throwable next = null) { 1346 super(operation, file, line, next); 1347 } 1348 1349 /++ 1350 The toString method will print out several components: 1351 1352 $(LIST 1353 * The file, line, static message, and object class name from the constructor. You can access these independently with the members `file`, `line`, `msg`, and [printableExceptionName]. 1354 * The generic category codes stored with this exception 1355 * Additional members stored with the exception child classes (e.g. platform error codes, associated function arguments) 1356 * The stack trace associated with the exception. You can access these lines independently with `foreach` over the `info` member. 1357 ) 1358 1359 This is meant to be read by the developer, not end users. You should wrap your user-relevant tasks in a try/catch block and construct more appropriate error messages from context available there, using the individual properties of the exception to add richness. 1360 +/ 1361 final override void toString(scope void delegate(in char[]) sink) const { 1362 // class name and info from constructor 1363 sink(printableExceptionName); 1364 sink("@"); 1365 sink(file); 1366 sink("("); 1367 char[16] buffer; 1368 sink(intToString(line, buffer[])); 1369 sink("): "); 1370 sink(message); 1371 1372 getAdditionalPrintableInformation((string name, in char[] value) { 1373 sink("\n"); 1374 sink(name); 1375 sink(": "); 1376 sink(value); 1377 }); 1378 1379 // full stack trace 1380 sink("\n----------------\n"); 1381 foreach(str; info) { 1382 sink(str); 1383 sink("\n"); 1384 } 1385 } 1386 /// ditto 1387 final override string toString() { 1388 string s; 1389 toString((in char[] chunk) { s ~= chunk; }); 1390 return s; 1391 } 1392 1393 /++ 1394 Users might like to see additional information with the exception. API consumers should pull this out of properties on your child class, but the parent class might not be able to deal with the arbitrary types at runtime the children can introduce, so bringing them all down to strings simplifies that. 1395 1396 Overrides should always call `super.getAdditionalPrintableInformation(sink);` before adding additional information by calling the sink with other arguments afterward. 1397 1398 You should spare no expense in preparing this information - translate error codes, build rich strings, whatever it takes - to make the information here useful to the reader. 1399 +/ 1400 void getAdditionalPrintableInformation(scope void delegate(string name, in char[] value) sink) const { 1401 1402 } 1403 1404 /++ 1405 This is the name of the exception class, suitable for printing. This should be static data (e.g. a string literal). Override it in subclasses. 1406 +/ 1407 string printableExceptionName() const { 1408 return typeid(this).name; 1409 } 1410 1411 /// deliberately hiding `Throwable.msg`. Use [message] and [toString] instead. 1412 @disable final void msg() {} 1413 1414 override const(char)[] message() const { 1415 return super.msg; 1416 } 1417 } 1418 1419 /++ 1420 1421 +/ 1422 class InvalidArgumentsException : ArsdExceptionBase { 1423 static struct InvalidArgument { 1424 string name; 1425 string description; 1426 LimitedVariant givenValue; 1427 } 1428 1429 InvalidArgument[] invalidArguments; 1430 1431 this(InvalidArgument[] invalidArguments, string functionName = __PRETTY_FUNCTION__, string file = __FILE__, size_t line = __LINE__, Throwable next = null) { 1432 this.invalidArguments = invalidArguments; 1433 super(functionName, file, line, next); 1434 } 1435 1436 this(string argumentName, string argumentDescription, LimitedVariant givenArgumentValue = LimitedVariant.init, string functionName = __PRETTY_FUNCTION__, string file = __FILE__, size_t line = __LINE__, Throwable next = null) { 1437 this([ 1438 InvalidArgument(argumentName, argumentDescription, givenArgumentValue) 1439 ], functionName, file, line, next); 1440 } 1441 1442 this(string argumentName, string argumentDescription, string functionName = __PRETTY_FUNCTION__, string file = __FILE__, size_t line = __LINE__, Throwable next = null) { 1443 this(argumentName, argumentDescription, LimitedVariant.init, functionName, file, line, next); 1444 } 1445 1446 override void getAdditionalPrintableInformation(scope void delegate(string name, in char[] value) sink) const { 1447 // FIXME: print the details better 1448 foreach(arg; invalidArguments) 1449 sink("invalidArguments[]", arg.name ~ " " ~ arg.description); 1450 } 1451 } 1452 1453 /++ 1454 Base class for when you've requested a feature that is not available. It may not be available because it is possible, but not yet implemented, or it might be because it is impossible on your operating system. 1455 +/ 1456 class FeatureUnavailableException : ArsdExceptionBase { 1457 this(string featureName = __PRETTY_FUNCTION__, string file = __FILE__, size_t line = __LINE__, Throwable next = null) { 1458 super(featureName, file, line, next); 1459 } 1460 } 1461 1462 /++ 1463 This means the feature could be done, but I haven't gotten around to implementing it yet. If you email me, I might be able to add it somewhat quickly and get back to you. 1464 +/ 1465 class NotYetImplementedException : FeatureUnavailableException { 1466 this(string featureName = __PRETTY_FUNCTION__, string file = __FILE__, size_t line = __LINE__, Throwable next = null) { 1467 super(featureName, file, line, next); 1468 } 1469 1470 } 1471 1472 /++ 1473 This means the feature is not supported by your current operating system. You might be able to get it in an update, but you might just have to find an alternate way of doing things. 1474 +/ 1475 class NotSupportedException : FeatureUnavailableException { 1476 this(string featureName, string file = __FILE__, size_t line = __LINE__, Throwable next = null) { 1477 super(featureName, file, line, next); 1478 } 1479 } 1480 1481 /++ 1482 This is a generic exception with attached arguments. It is used when I had to throw something but didn't want to write a new class. 1483 1484 You can catch an ArsdException to get its passed arguments out. 1485 1486 You can pass either a base class or a string as `Type`. 1487 1488 See the examples for how to use it. 1489 +/ 1490 template ArsdException(alias Type, DataTuple...) { 1491 static if(DataTuple.length) 1492 alias Parent = ArsdException!(Type, DataTuple[0 .. $-1]); 1493 else 1494 alias Parent = ArsdExceptionBase; 1495 1496 class ArsdException : Parent { 1497 DataTuple data; 1498 1499 this(DataTuple data, string file = __FILE__, size_t line = __LINE__) { 1500 this.data = data; 1501 static if(is(Parent == ArsdExceptionBase)) 1502 super(null, file, line); 1503 else 1504 super(data[0 .. $-1], file, line); 1505 } 1506 1507 static opCall(R...)(R r, string file = __FILE__, size_t line = __LINE__) { 1508 return new ArsdException!(Type, DataTuple, R)(r, file, line); 1509 } 1510 1511 override string printableExceptionName() const { 1512 static if(DataTuple.length) 1513 enum str = "ArsdException!(" ~ Type.stringof ~ ", " ~ DataTuple.stringof[1 .. $-1] ~ ")"; 1514 else 1515 enum str = "ArsdException!" ~ Type.stringof; 1516 return str; 1517 } 1518 1519 override void getAdditionalPrintableInformation(scope void delegate(string name, in char[] value) sink) const { 1520 ArsdExceptionBase.getAdditionalPrintableInformation(sink); 1521 1522 foreach(idx, datum; data) { 1523 enum int lol = cast(int) idx; 1524 enum key = "[" ~ lol.stringof ~ "] " ~ DataTuple[idx].stringof; 1525 sink(key, toStringInternal(datum)); 1526 } 1527 } 1528 } 1529 } 1530 1531 /// This example shows how you can throw and catch the ad-hoc exception types. 1532 unittest { 1533 // you can throw and catch by matching the string and argument types 1534 try { 1535 // throw it with parenthesis after the template args (it uses opCall to construct) 1536 throw ArsdException!"Test"(); 1537 // you could also `throw new ArsdException!"test";`, but that gets harder with args 1538 // as we'll see in the following example 1539 assert(0); // remove from docs 1540 } catch(ArsdException!"Test" e) { // catch it without them 1541 // this has no useful information except for the type 1542 // but you can catch it like this and it is still more than generic Exception 1543 } 1544 1545 // an exception's job is to deliver useful information up the chain 1546 // and you can do that easily by passing arguments: 1547 1548 try { 1549 throw ArsdException!"Test"(4, "four"); 1550 // you could also `throw new ArsdException!("Test", int, string)(4, "four")` 1551 // but now you start to see how the opCall convenience constructor simplifies things 1552 assert(0); // remove from docs 1553 } catch(ArsdException!("Test", int, string) e) { // catch it and use info by specifying types 1554 assert(e.data[0] == 4); // and extract arguments like this 1555 assert(e.data[1] == "four"); 1556 } 1557 1558 // a throw site can add additional information without breaking code that catches just some 1559 // generally speaking, each additional argument creates a new subclass on top of the previous args 1560 // so you can cast 1561 1562 try { 1563 throw ArsdException!"Test"(4, "four", 9); 1564 assert(0); // remove from docs 1565 } catch(ArsdException!("Test", int, string) e) { // this catch still works 1566 assert(e.data[0] == 4); 1567 assert(e.data[1] == "four"); 1568 // but if you were to print it, all the members would be there 1569 // import std.stdio; writeln(e); // would show something like: 1570 /+ 1571 ArsdException!("Test", int, string, int)@file.d(line): 1572 [0] int: 4 1573 [1] string: four 1574 [2] int: 9 1575 +/ 1576 // indicating that there's additional information available if you wanted to process it 1577 1578 // and meanwhile: 1579 ArsdException!("Test", int) e2 = e; // this implicit cast works thanks to the parent-child relationship 1580 ArsdException!"Test" e3 = e; // this works too, the base type/string still matches 1581 1582 // so catching those types would work too 1583 } 1584 } 1585 1586 /++ 1587 A tagged union that holds an error code from system apis, meaning one from Windows GetLastError() or C's errno. 1588 1589 You construct it with `SystemErrorCode(thing)` and the overloaded constructor tags and stores it. 1590 +/ 1591 struct SystemErrorCode { 1592 /// 1593 enum Type { 1594 errno, /// 1595 win32 /// 1596 } 1597 1598 const Type type; /// 1599 const int code; /// You should technically cast it back to DWORD if it is a win32 code 1600 1601 /++ 1602 C/unix error are typed as signed ints... 1603 Windows' errors are typed DWORD, aka unsigned... 1604 1605 so just passing them straight up will pick the right overload here to set the tag. 1606 +/ 1607 this(int errno) { 1608 this.type = Type.errno; 1609 this.code = errno; 1610 } 1611 1612 /// ditto 1613 this(uint win32) { 1614 this.type = Type.win32; 1615 this.code = win32; 1616 } 1617 1618 /++ 1619 Returns if the code indicated success. 1620 1621 Please note that many calls do not actually set a code to success, but rather just don't touch it. Thus this may only be true on `init`. 1622 +/ 1623 bool wasSuccessful() const { 1624 final switch(type) { 1625 case Type.errno: 1626 return this.code == 0; 1627 case Type.win32: 1628 return this.code == 0; 1629 } 1630 } 1631 1632 /++ 1633 Constructs a string containing both the code and the explanation string. 1634 +/ 1635 string toString() const { 1636 return "[" ~ codeAsString ~ "] " ~ errorString; 1637 } 1638 1639 /++ 1640 The numeric code itself as a string. 1641 1642 See [errorString] for a text explanation of the code. 1643 +/ 1644 string codeAsString() const { 1645 char[16] buffer; 1646 final switch(type) { 1647 case Type.errno: 1648 return intToString(code, buffer[]).idup; 1649 case Type.win32: 1650 buffer[0 .. 2] = "0x"; 1651 return buffer[0 .. 2 + intToString(cast(uint) code, buffer[2 .. $], IntToStringArgs().withRadix(16).withPadding(8)).length].idup; 1652 } 1653 } 1654 1655 /++ 1656 A text explanation of the code. See [codeAsString] for a string representation of the numeric representation. 1657 +/ 1658 string errorString() const { 1659 final switch(type) { 1660 case Type.errno: 1661 import core.stdc.string; 1662 auto strptr = strerror(code); 1663 auto orig = strptr; 1664 int len; 1665 while(*strptr++) { 1666 len++; 1667 } 1668 1669 return orig[0 .. len].idup; 1670 case Type.win32: 1671 version(Windows) { 1672 wchar[256] buffer; 1673 auto size = FormatMessageW( 1674 FORMAT_MESSAGE_FROM_SYSTEM | FORMAT_MESSAGE_IGNORE_INSERTS, 1675 null, 1676 code, 1677 MAKELANGID(LANG_NEUTRAL, SUBLANG_DEFAULT), 1678 buffer.ptr, 1679 buffer.length, 1680 null 1681 ); 1682 1683 return makeUtf8StringFromWindowsString(buffer[0 .. size]).stripInternal; 1684 } else { 1685 return null; 1686 } 1687 } 1688 } 1689 } 1690 1691 /++ 1692 1693 +/ 1694 struct SavedArgument { 1695 string name; 1696 LimitedVariant value; 1697 } 1698 1699 /++ 1700 1701 +/ 1702 class SystemApiException : ArsdExceptionBase { 1703 this(string msg, int originalErrorNo, scope SavedArgument[] args = null, string file = __FILE__, size_t line = __LINE__, Throwable next = null) { 1704 this(msg, SystemErrorCode(originalErrorNo), args, file, line, next); 1705 } 1706 1707 version(Windows) 1708 this(string msg, DWORD windowsError, scope SavedArgument[] args = null, string file = __FILE__, size_t line = __LINE__, Throwable next = null) { 1709 this(msg, SystemErrorCode(windowsError), args, file, line, next); 1710 } 1711 1712 this(string msg, SystemErrorCode code, SavedArgument[] args = null, string file = __FILE__, size_t line = __LINE__, Throwable next = null) { 1713 this.errorCode = code; 1714 1715 // discard stuff that won't fit 1716 if(args.length > this.args.length) 1717 args = args[0 .. this.args.length]; 1718 1719 this.args[0 .. args.length] = args[]; 1720 1721 super(msg, file, line, next); 1722 } 1723 1724 /++ 1725 1726 +/ 1727 const SystemErrorCode errorCode; 1728 1729 /++ 1730 1731 +/ 1732 const SavedArgument[8] args; 1733 1734 override void getAdditionalPrintableInformation(scope void delegate(string name, in char[] value) sink) const { 1735 super.getAdditionalPrintableInformation(sink); 1736 sink("Error code", errorCode.toString()); 1737 1738 foreach(arg; args) 1739 if(arg.name !is null) 1740 sink(arg.name, arg.value.toString()); 1741 } 1742 1743 } 1744 1745 /++ 1746 The low level use of this would look like `throw new WindowsApiException("MsgWaitForMultipleObjectsEx", GetLastError())` but it is meant to be used from higher level things like [Win32Enforce]. 1747 1748 History: 1749 Moved from simpledisplay.d to core.d in March 2023 (dub v11.0). 1750 +/ 1751 alias WindowsApiException = SystemApiException; 1752 1753 /++ 1754 History: 1755 Moved from simpledisplay.d to core.d in March 2023 (dub v11.0). 1756 +/ 1757 alias ErrnoApiException = SystemApiException; 1758 1759 /++ 1760 Calls the C API function `fn`. If it returns an error value, it throws an [ErrnoApiException] (or subclass) after getting `errno`. 1761 +/ 1762 template ErrnoEnforce(alias fn, alias errorValue = void) { 1763 static if(is(typeof(fn) Return == return)) 1764 static if(is(typeof(fn) Params == __parameters)) { 1765 static if(is(errorValue == void)) { 1766 static if(is(typeof(null) : Return)) 1767 enum errorValueToUse = null; 1768 else static if(is(Return : long)) 1769 enum errorValueToUse = -1; 1770 else 1771 static assert(0, "Please pass the error value"); 1772 } else { 1773 enum errorValueToUse = errorValue; 1774 } 1775 1776 Return ErrnoEnforce(Params params, ArgSentinel sentinel = ArgSentinel.init, string file = __FILE__, size_t line = __LINE__) { 1777 import core.stdc.errno; 1778 1779 Return value = fn(params); 1780 1781 if(value == errorValueToUse) { 1782 SavedArgument[] args; // FIXME 1783 /+ 1784 static foreach(idx; 0 .. Params.length) 1785 args ~= SavedArgument( 1786 __traits(identifier, Params[idx .. idx + 1]), 1787 params[idx] 1788 ); 1789 +/ 1790 throw new ErrnoApiException(__traits(identifier, fn), errno, args, file, line); 1791 } 1792 1793 return value; 1794 } 1795 } 1796 } 1797 1798 version(Windows) { 1799 /++ 1800 Calls the Windows API function `fn`. If it returns an error value, it throws a [WindowsApiException] (or subclass) after calling `GetLastError()`. 1801 +/ 1802 template Win32Enforce(alias fn, alias errorValue = void) { 1803 static if(is(typeof(fn) Return == return)) 1804 static if(is(typeof(fn) Params == __parameters)) { 1805 static if(is(errorValue == void)) { 1806 static if(is(Return == BOOL)) 1807 enum errorValueToUse = false; 1808 else static if(is(Return : HANDLE)) 1809 enum errorValueToUse = NULL; 1810 else static if(is(Return == DWORD)) 1811 enum errorValueToUse = cast(DWORD) 0xffffffff; 1812 else 1813 static assert(0, "Please pass the error value"); 1814 } else { 1815 enum errorValueToUse = errorValue; 1816 } 1817 1818 Return Win32Enforce(Params params, ArgSentinel sentinel = ArgSentinel.init, string file = __FILE__, size_t line = __LINE__) { 1819 Return value = fn(params); 1820 1821 if(value == errorValueToUse) { 1822 auto error = GetLastError(); 1823 SavedArgument[] args; // FIXME 1824 throw new WindowsApiException(__traits(identifier, fn), error, args, file, line); 1825 } 1826 1827 return value; 1828 } 1829 } 1830 } 1831 1832 } 1833 1834 /+ 1835 =============== 1836 EVENT LOOP CORE 1837 =============== 1838 +/ 1839 1840 /+ 1841 UI threads 1842 need to get window messages in addition to all the other jobs 1843 I/O Worker threads 1844 need to get commands for read/writes, run them, and send the reply back. not necessary on Windows 1845 if interrupted, check cancel flags. 1846 CPU Worker threads 1847 gets functions, runs them, send reply back. should send a cancel flag to periodically check 1848 Task worker threads 1849 runs fibers and multiplexes them 1850 1851 1852 General procedure: 1853 issue the read/write command 1854 if it would block on linux, epoll associate it. otherwise do the callback immediately 1855 1856 callbacks have default affinity to the current thread, meaning their callbacks always run here 1857 accepts can usually be dispatched to any available thread tho 1858 1859 // In other words, a single thread can be associated with, at most, one I/O completion port. 1860 1861 Realistically, IOCP only used if there is no thread affinity. If there is, just do overlapped w/ sleepex. 1862 1863 1864 case study: http server 1865 1866 1) main thread starts the server. it does an accept loop with no thread affinity. the main thread does NOT check the global queue (the iocp/global epoll) 1867 2) connections come in and are assigned to first available thread via the iocp/global epoll 1868 3) these run local event loops until the connection task is finished 1869 1870 EVENT LOOP TYPES: 1871 1) main ui thread - MsgWaitForMultipleObjectsEx / epoll on the local ui. it does NOT check the any worker thread thing! 1872 The main ui thread should never terminate until the program is ready to close. 1873 You can have additional ui threads in theory but im not really gonna support that in full; most things will assume there is just the one. simpledisplay's gui thread is the primary if it exists. (and sdpy will prolly continue to be threaded the way it is now) 1874 1875 The biggest complication is the TerminalDirectToEmulator, where the primary ui thread is NOT the thread that runs `main` 1876 2) worker thread GetQueuedCompletionStatusEx / epoll on the local thread fd and the global epoll fd 1877 3) local event loop - check local things only. SleepEx / epoll on local thread fd. This more of a compatibility hack for `waitForCompletion` outside a fiber. 1878 1879 i'll use: 1880 * QueueUserAPC to send interruptions to a worker thread 1881 * PostQueuedCompletionStatus is to send interruptions to any available thread. 1882 * PostMessage to a window 1883 * ??? to a fiber task 1884 1885 I also need a way to de-duplicate events in the queue so if you try to push the same thing it won't trigger multiple times.... I might want to keep a duplicate of the thing... really, what I'd do is post the "event wake up" message and keep the queue in my own thing. (WM_PAINT auto-coalesces) 1886 1887 Destructors need to be able to post messages back to a specific task to queue thread-affinity cleanup. This must be GC safe. 1888 1889 A task might want to wait on certain events. If the task is a fiber, it yields and gets called upon the event. If the task is a thread, it really has to call the event loop... which can be a loop of loops we want to avoid. `waitForCompletion` is more often gonna be used just to run the loop at top level tho... it might not even check for the global info availability so it'd run the local thing only. 1890 1891 APCs should not themselves enter an alterable wait cuz it can stack overflow. So generally speaking, they should avoid calling fibers or other event loops. 1892 +/ 1893 1894 /++ 1895 You can also pass a handle to a specific thread, if you have one. 1896 +/ 1897 enum ThreadToRunIn { 1898 /++ 1899 The callback should be only run by the same thread that set it. 1900 +/ 1901 CurrentThread, 1902 /++ 1903 The UI thread is a special one - it is the supervisor of the workers and the controller of gui and console handles. It is the first thread to call [arsd_core_init] actively running an event loop unless there is a thread that has actively asserted the ui supervisor role. FIXME is this true after i implemen it? 1904 1905 A ui thread should be always quickly responsive to new events. 1906 1907 There should only be one main ui thread, in which simpledisplay and minigui can be used. 1908 1909 Other threads can run like ui threads, but are considered temporary and only concerned with their own needs (it is the default style of loop 1910 for an undeclared thread but will not receive messages from other threads unless there is no other option) 1911 1912 1913 Ad-Hoc thread - something running an event loop that isn't another thing 1914 Controller thread - running an explicit event loop instance set as not a task runner or blocking worker 1915 UI thread - simpledisplay's event loop, which it will require remain live for the duration of the program (running two .eventLoops without a parent EventLoop instance will become illegal, throwing at runtime if it happens telling people to change their code) 1916 1917 Windows HANDLES will always be listened on the thread itself that is requesting, UNLESS it is a worker/helper thread, in which case it goes to a coordinator thread. since it prolly can't rely on the parent per se this will have to be one created by arsd core init, UNLESS the parent is inside an explicit EventLoop structure. 1918 1919 All use the MsgWaitForMultipleObjectsEx pattern 1920 1921 1922 +/ 1923 UiThread, 1924 /++ 1925 The callback can be called from any available worker thread. It will be added to a global queue and the first thread to see it will run it. 1926 1927 These will not run on the UI thread unless there is no other option on the platform (and all platforms this lib supports have other options). 1928 1929 These are expected to run cooperatively multitasked things; functions that frequently yield as they wait on other tasks. Think a fiber. 1930 1931 A task runner should be generally responsive to new events. 1932 +/ 1933 AnyAvailableTaskRunnerThread, 1934 /++ 1935 These are expected to run longer blocking, but independent operations. Think an individual function with no context. 1936 1937 A blocking worker can wait hundreds of milliseconds between checking for new events. 1938 +/ 1939 AnyAvailableBlockingWorkerThread, 1940 /++ 1941 The callback will be duplicated across all threads known to the arsd.core event loop. 1942 1943 It adds it to an immutable queue that each thread will go through... might just replace with an exit() function. 1944 1945 1946 so to cancel all associated tasks for like a web server, it could just have the tasks atomicAdd to a counter and subtract when they are finished. Then you have a single semaphore you signal the number of times you have an active thing and wait for them to acknowledge it. 1947 1948 threads should report when they start running the loop and they really should report when they terminate but that isn't reliable 1949 1950 1951 hmmm what if: all user-created threads (the public api) count as ui threads. only ones created in here are task runners or helpers. ui threads can wait on a global event to exit. 1952 1953 there's still prolly be one "the" ui thread, which does the handle listening on windows and is the one sdpy wants. 1954 +/ 1955 BroadcastToAllThreads, 1956 } 1957 1958 /++ 1959 Initializes the arsd core event loop and creates its worker threads. You don't actually have to call this, since the first use of an arsd.core function that requires it will call it implicitly, but calling it yourself gives you a chance to control the configuration more explicitly if you want to. 1960 +/ 1961 void arsd_core_init(int numberOfWorkers = 0) { 1962 1963 } 1964 1965 version(Windows) 1966 class WindowsHandleReader_ex { 1967 // Windows handles are always dispatched to the main ui thread, which can then send a command back to a worker thread to run the callback if needed 1968 this(HANDLE handle) {} 1969 } 1970 1971 version(Posix) 1972 class PosixFdReader_ex { 1973 // posix readers can just register with whatever instance we want to handle the callback 1974 } 1975 1976 /++ 1977 1978 +/ 1979 interface ICoreEventLoop { 1980 /++ 1981 Runs the event loop for this thread until the `until` delegate returns `true`. 1982 +/ 1983 final void run(scope bool delegate() until) { 1984 while(!until()) { 1985 runOnce(); 1986 } 1987 } 1988 1989 /++ 1990 Runs a single iteration of the event loop for this thread. It will return when the first thing happens, but that thing might be totally uninteresting to anyone, or it might trigger significant work you'll wait on. 1991 +/ 1992 void runOnce(); 1993 1994 // to send messages between threads, i'll queue up a function that just call dispatchMessage. can embed the arg inside the callback helper prolly. 1995 // tho i might prefer to actually do messages w/ run payloads so it is easier to deduplicate i can still dedupe by insepcting the call args so idk 1996 1997 version(Posix) { 1998 @mustuse 1999 static struct UnregisterToken { 2000 private CoreEventLoopImplementation impl; 2001 private int fd; 2002 private CallbackHelper cb; 2003 2004 /++ 2005 Unregisters the file descriptor from the event loop and releases the reference to the callback held by the event loop (which will probably free it). 2006 2007 You must call this when you're done. Normally, this will be right before you close the fd (Which is often after the other side closes it, meaning you got a 0 length read). 2008 +/ 2009 void unregister() { 2010 assert(impl !is null, "Cannot reuse unregister token"); 2011 2012 version(Arsd_core_epoll) { 2013 impl.unregisterFd(fd); 2014 } else version(Arsd_core_kqueue) { 2015 // intentionally blank - all registrations are one-shot there 2016 // FIXME: actually it might not have gone off yet, in that case we do need to delete the filter 2017 } else version(EmptyCoreEvent) { 2018 2019 } 2020 else static assert(0); 2021 2022 cb.release(); 2023 this = typeof(this).init; 2024 } 2025 } 2026 2027 @mustuse 2028 static struct RearmToken { 2029 private bool readable; 2030 private CoreEventLoopImplementation impl; 2031 private int fd; 2032 private CallbackHelper cb; 2033 private uint flags; 2034 2035 /++ 2036 Calls [UnregisterToken.unregister] 2037 +/ 2038 void unregister() { 2039 assert(impl !is null, "cannot reuse rearm token after unregistering it"); 2040 2041 version(Arsd_core_epoll) { 2042 impl.unregisterFd(fd); 2043 } else version(Arsd_core_kqueue) { 2044 // intentionally blank - all registrations are one-shot there 2045 // FIXME: actually it might not have gone off yet, in that case we do need to delete the filter 2046 } else version(EmptyCoreEvent) { 2047 2048 } else static assert(0); 2049 2050 cb.release(); 2051 this = typeof(this).init; 2052 } 2053 2054 /++ 2055 Rearms the event so you will get another callback next time it is ready. 2056 +/ 2057 void rearm() { 2058 assert(impl !is null, "cannot reuse rearm token after unregistering it"); 2059 impl.rearmFd(this); 2060 } 2061 } 2062 2063 UnregisterToken addCallbackOnFdReadable(int fd, CallbackHelper cb); 2064 RearmToken addCallbackOnFdReadableOneShot(int fd, CallbackHelper cb); 2065 RearmToken addCallbackOnFdWritableOneShot(int fd, CallbackHelper cb); 2066 } 2067 } 2068 2069 /++ 2070 Get the event loop associated with this thread 2071 +/ 2072 ICoreEventLoop getThisThreadEventLoop(EventLoopType type = EventLoopType.AdHoc) { 2073 static ICoreEventLoop loop; 2074 if(loop is null) 2075 loop = new CoreEventLoopImplementation(); 2076 return loop; 2077 } 2078 2079 /++ 2080 The internal types that will be exposed through other api things. 2081 +/ 2082 package(arsd) enum EventLoopType { 2083 /++ 2084 The event loop is being run temporarily and the thread doesn't promise to keep running it. 2085 +/ 2086 AdHoc, 2087 /++ 2088 The event loop struct has been instantiated at top level. Its destructor will run when the 2089 function exits, which is only at the end of the entire block of work it is responsible for. 2090 2091 It must be in scope for the whole time the arsd event loop functions are expected to be used 2092 (meaning it should generally be top-level in `main`) 2093 +/ 2094 Explicit, 2095 /++ 2096 A specialization of `Explicit`, so all the same rules apply there, but this is specifically the event loop coming from simpledisplay or minigui. It will run for the duration of the UI's existence. 2097 +/ 2098 Ui, 2099 /++ 2100 A special event loop specifically for threads that listen to the task runner queue and handle I/O events from running tasks. Typically, a task runner runs cooperatively multitasked coroutines (so they prefer not to block the whole thread). 2101 +/ 2102 TaskRunner, 2103 /++ 2104 A special event loop specifically for threads that listen to the helper function request queue. Helper functions are expected to run independently for a somewhat long time (them blocking the thread for some time is normal) and send a reply message back to the requester. 2105 +/ 2106 HelperWorker 2107 } 2108 2109 /+ 2110 Tasks are given an object to talk to their parent... can be a dialog where it is like 2111 2112 sendBuffer 2113 waitForWordToProceed 2114 2115 in a loop 2116 2117 2118 Tasks are assigned to a worker thread and may share it with other tasks. 2119 +/ 2120 2121 2122 // the GC may not be able to see this! remember, it can be hidden inside kernel buffers 2123 version(HasThread) private class CallbackHelper { 2124 import core.memory; 2125 2126 void call() { 2127 if(callback) 2128 callback(); 2129 } 2130 2131 void delegate() callback; 2132 void*[3] argsStore; 2133 2134 void addref() { 2135 atomicOp!"+="(refcount, 1); 2136 } 2137 2138 void release() { 2139 if(atomicOp!"-="(refcount, 1) <= 0) { 2140 if(flags & 1) 2141 GC.removeRoot(cast(void*) this); 2142 } 2143 } 2144 2145 private shared(int) refcount; 2146 private uint flags; 2147 2148 this(void function() callback) { 2149 this( () { callback(); } ); 2150 } 2151 2152 this(void delegate() callback, bool addRoot = true) { 2153 if(addRoot) { 2154 GC.addRoot(cast(void*) this); 2155 this.flags |= 1; 2156 } 2157 2158 this.addref(); 2159 this.callback = callback; 2160 } 2161 } 2162 2163 /++ 2164 This represents a file. Technically, file paths aren't actually strings (for example, on Linux, they need not be valid utf-8, while a D string is supposed to be), even though we almost always use them like that. 2165 2166 This type is meant to represent a filename / path. I might not keep it around. 2167 +/ 2168 struct FilePath { 2169 string path; 2170 2171 bool isNull() { 2172 return path is null; 2173 } 2174 2175 bool opCast(T:bool)() { 2176 return !isNull; 2177 } 2178 2179 string toString() { 2180 return path; 2181 } 2182 2183 //alias toString this; 2184 } 2185 2186 /++ 2187 Represents a generic async, waitable request. 2188 +/ 2189 class AsyncOperationRequest { 2190 /++ 2191 Actually issues the request, starting the operation. 2192 +/ 2193 abstract void start(); 2194 /++ 2195 Cancels the request. This will cause `isComplete` to return true once the cancellation has been processed, but [AsyncOperationResponse.wasSuccessful] will return `false` (unless it completed before the cancellation was processed, in which case it is still allowed to finish successfully). 2196 2197 After cancelling a request, you should still wait for it to complete to ensure that the task has actually released its resources before doing anything else on it. 2198 2199 Once a cancellation request has been sent, it cannot be undone. 2200 +/ 2201 abstract void cancel(); 2202 2203 /++ 2204 Returns `true` if the operation has been completed. It may be completed successfully, cancelled, or have errored out - to check this, call [waitForCompletion] and check the members on the response object. 2205 +/ 2206 abstract bool isComplete(); 2207 /++ 2208 Waits until the request has completed - successfully or otherwise - and returns the response object. It will run an ad-hoc event loop that may call other callbacks while waiting. 2209 2210 The response object may be embedded in the request object - do not reuse the request until you are finished with the response and do not keep the response around longer than you keep the request. 2211 2212 2213 Note to implementers: all subclasses should override this and return their specific response object. You can use the top-level `waitForFirstToCompleteByIndex` function with a single-element static array to help with the implementation. 2214 +/ 2215 abstract AsyncOperationResponse waitForCompletion(); 2216 2217 /++ 2218 2219 +/ 2220 // abstract void repeat(); 2221 } 2222 2223 /++ 2224 2225 +/ 2226 interface AsyncOperationResponse { 2227 /++ 2228 Returns true if the request completed successfully, finishing what it was supposed to. 2229 2230 Should be set to `false` if the request was cancelled before completing or encountered an error. 2231 +/ 2232 bool wasSuccessful(); 2233 } 2234 2235 /++ 2236 It returns the $(I request) so you can identify it more easily. `request.waitForCompletion()` is guaranteed to return the response without any actual wait, since it is already complete when this function returns. 2237 2238 Please note that "completion" is not necessary successful completion; a request being cancelled or encountering an error also counts as it being completed. 2239 2240 The `waitForFirstToCompleteByIndex` version instead returns the index of the array entry that completed first. 2241 2242 It is your responsibility to remove the completed request from the array before calling the function again, since any request already completed will always be immediately returned. 2243 2244 You might prefer using [asTheyComplete], which will give each request as it completes and loop over until all of them are complete. 2245 2246 Returns: 2247 `null` or `requests.length` if none completed before returning. 2248 +/ 2249 AsyncOperationRequest waitForFirstToComplete(AsyncOperationRequest[] requests...) { 2250 auto idx = waitForFirstToCompleteByIndex(requests); 2251 if(idx == requests.length) 2252 return null; 2253 return requests[idx]; 2254 } 2255 /// ditto 2256 size_t waitForFirstToCompleteByIndex(AsyncOperationRequest[] requests...) { 2257 size_t helper() { 2258 foreach(idx, request; requests) 2259 if(request.isComplete()) 2260 return idx; 2261 return requests.length; 2262 } 2263 2264 auto idx = helper(); 2265 // if one is already done, return it 2266 if(idx != requests.length) 2267 return idx; 2268 2269 // otherwise, run the ad-hoc event loop until one is 2270 // FIXME: what if we are inside a fiber? 2271 auto el = getThisThreadEventLoop(); 2272 el.run(() => (idx = helper()) != requests.length); 2273 2274 return idx; 2275 } 2276 2277 /++ 2278 Waits for all the `requests` to complete, giving each one through the range interface as it completes. 2279 2280 This meant to be used in a foreach loop. 2281 2282 The `requests` array and its contents must remain valid for the lifetime of the returned range. Its contents may be shuffled as the requests complete (the implementation works through an unstable sort+remove). 2283 +/ 2284 AsTheyCompleteRange asTheyComplete(AsyncOperationRequest[] requests...) { 2285 return AsTheyCompleteRange(requests); 2286 } 2287 /// ditto 2288 struct AsTheyCompleteRange { 2289 AsyncOperationRequest[] requests; 2290 2291 this(AsyncOperationRequest[] requests) { 2292 this.requests = requests; 2293 2294 if(requests.length == 0) 2295 return; 2296 2297 // wait for first one to complete, then move it to the front of the array 2298 moveFirstCompleteToFront(); 2299 } 2300 2301 private void moveFirstCompleteToFront() { 2302 auto idx = waitForFirstToCompleteByIndex(requests); 2303 2304 auto tmp = requests[0]; 2305 requests[0] = requests[idx]; 2306 requests[idx] = tmp; 2307 } 2308 2309 bool empty() { 2310 return requests.length == 0; 2311 } 2312 2313 void popFront() { 2314 assert(!empty); 2315 /+ 2316 this needs to 2317 1) remove the front of the array as being already processed (unless it is the initial priming call) 2318 2) wait for one of them to complete 2319 3) move the complete one to the front of the array 2320 +/ 2321 2322 requests[0] = requests[$-1]; 2323 requests = requests[0 .. $-1]; 2324 2325 if(requests.length) 2326 moveFirstCompleteToFront(); 2327 } 2328 2329 AsyncOperationRequest front() { 2330 return requests[0]; 2331 } 2332 } 2333 2334 version(Windows) { 2335 alias NativeFileHandle = HANDLE; /// 2336 alias NativeSocketHandle = SOCKET; /// 2337 alias NativePipeHandle = HANDLE; /// 2338 } else version(Posix) { 2339 alias NativeFileHandle = int; /// 2340 alias NativeSocketHandle = int; /// 2341 alias NativePipeHandle = int; /// 2342 } 2343 2344 /++ 2345 An `AbstractFile` represents a file handle on the operating system level. You cannot do much with it. 2346 +/ 2347 version(HasFile) class AbstractFile { 2348 private { 2349 NativeFileHandle handle; 2350 } 2351 2352 /++ 2353 +/ 2354 enum OpenMode { 2355 readOnly, /// C's "r", the file is read 2356 writeWithTruncation, /// C's "w", the file is blanked upon opening so it only holds what you write 2357 appendOnly, /// C's "a", writes will always be appended to the file 2358 readAndWrite /// C's "r+", writes will overwrite existing parts of the file based on where you seek (default is at the beginning) 2359 } 2360 2361 /++ 2362 +/ 2363 enum RequirePreexisting { 2364 no, 2365 yes 2366 } 2367 2368 /+ 2369 enum SpecialFlags { 2370 randomAccessExpected, /// FILE_FLAG_SEQUENTIAL_SCAN is turned off and posix_fadvise(POSIX_FADV_SEQUENTIAL) 2371 skipCache, /// O_DSYNC, FILE_FLAG_NO_BUFFERING and maybe WRITE_THROUGH. note that metadata still goes through the cache, FlushFileBuffers and fsync can still do those 2372 temporary, /// FILE_ATTRIBUTE_TEMPORARY on Windows, idk how to specify on linux. also FILE_FLAG_DELETE_ON_CLOSE can be combined to make a (almost) all memory file. kinda like a private anonymous mmap i believe. 2373 deleteWhenClosed, /// Windows has a flag for this but idk if it is of any real use 2374 async, /// open it in overlapped mode, all reads and writes must then provide an offset. Only implemented on Windows 2375 } 2376 +/ 2377 2378 /++ 2379 2380 +/ 2381 protected this(bool async, FilePath filename, OpenMode mode = OpenMode.readOnly, RequirePreexisting require = RequirePreexisting.no, uint specialFlags = 0) { 2382 version(Windows) { 2383 DWORD access; 2384 DWORD creation; 2385 2386 final switch(mode) { 2387 case OpenMode.readOnly: 2388 access = GENERIC_READ; 2389 creation = OPEN_EXISTING; 2390 break; 2391 case OpenMode.writeWithTruncation: 2392 access = GENERIC_WRITE; 2393 2394 final switch(require) { 2395 case RequirePreexisting.no: 2396 creation = CREATE_ALWAYS; 2397 break; 2398 case RequirePreexisting.yes: 2399 creation = TRUNCATE_EXISTING; 2400 break; 2401 } 2402 break; 2403 case OpenMode.appendOnly: 2404 access = FILE_APPEND_DATA; 2405 2406 final switch(require) { 2407 case RequirePreexisting.no: 2408 creation = CREATE_ALWAYS; 2409 break; 2410 case RequirePreexisting.yes: 2411 creation = OPEN_EXISTING; 2412 break; 2413 } 2414 break; 2415 case OpenMode.readAndWrite: 2416 access = GENERIC_READ | GENERIC_WRITE; 2417 2418 final switch(require) { 2419 case RequirePreexisting.no: 2420 creation = CREATE_NEW; 2421 break; 2422 case RequirePreexisting.yes: 2423 creation = OPEN_EXISTING; 2424 break; 2425 } 2426 break; 2427 } 2428 2429 WCharzBuffer wname = WCharzBuffer(filename.path); 2430 2431 auto handle = CreateFileW( 2432 wname.ptr, 2433 access, 2434 FILE_SHARE_READ, 2435 null, 2436 creation, 2437 FILE_ATTRIBUTE_NORMAL | (async ? FILE_FLAG_OVERLAPPED : 0), 2438 null 2439 ); 2440 2441 if(handle == INVALID_HANDLE_VALUE) { 2442 // FIXME: throw the filename and other params here too 2443 SavedArgument[3] args; 2444 args[0] = SavedArgument("filename", LimitedVariant(filename.path)); 2445 args[1] = SavedArgument("access", LimitedVariant(access, 2)); 2446 args[2] = SavedArgument("requirePreexisting", LimitedVariant(require == RequirePreexisting.yes)); 2447 throw new WindowsApiException("CreateFileW", GetLastError(), args[]); 2448 } 2449 2450 this.handle = handle; 2451 } else version(Posix) { 2452 import core.sys.posix.unistd; 2453 import core.sys.posix.fcntl; 2454 2455 CharzBuffer namez = CharzBuffer(filename.path); 2456 int flags; 2457 2458 // FIXME does mac not have cloexec for real or is this just a druntime problem????? 2459 version(Arsd_core_has_cloexec) { 2460 flags = O_CLOEXEC; 2461 } else { 2462 scope(success) 2463 setCloExec(this.handle); 2464 } 2465 2466 if(async) 2467 flags |= O_NONBLOCK; 2468 2469 final switch(mode) { 2470 case OpenMode.readOnly: 2471 flags |= O_RDONLY; 2472 break; 2473 case OpenMode.writeWithTruncation: 2474 flags |= O_WRONLY | O_TRUNC; 2475 2476 final switch(require) { 2477 case RequirePreexisting.no: 2478 flags |= O_CREAT; 2479 break; 2480 case RequirePreexisting.yes: 2481 break; 2482 } 2483 break; 2484 case OpenMode.appendOnly: 2485 flags |= O_APPEND; 2486 2487 final switch(require) { 2488 case RequirePreexisting.no: 2489 flags |= O_CREAT; 2490 break; 2491 case RequirePreexisting.yes: 2492 break; 2493 } 2494 break; 2495 case OpenMode.readAndWrite: 2496 flags |= O_RDWR; 2497 2498 final switch(require) { 2499 case RequirePreexisting.no: 2500 flags |= O_CREAT; 2501 break; 2502 case RequirePreexisting.yes: 2503 break; 2504 } 2505 break; 2506 } 2507 2508 auto perms = S_IRUSR | S_IWUSR | S_IRGRP | S_IROTH; 2509 int fd = open(namez.ptr, flags, perms); 2510 if(fd == -1) { 2511 SavedArgument[3] args; 2512 args[0] = SavedArgument("filename", LimitedVariant(filename.path)); 2513 args[1] = SavedArgument("flags", LimitedVariant(flags, 2)); 2514 args[2] = SavedArgument("perms", LimitedVariant(perms, 8)); 2515 throw new ErrnoApiException("open", errno, args[]); 2516 } 2517 2518 this.handle = fd; 2519 } 2520 } 2521 2522 /++ 2523 2524 +/ 2525 private this(NativeFileHandle handleToWrap) { 2526 this.handle = handleToWrap; 2527 } 2528 2529 // only available on some types of file 2530 long size() { return 0; } 2531 2532 // note that there is no fsync thing, instead use the special flag. 2533 2534 /++ 2535 2536 +/ 2537 void close() { 2538 version(Windows) { 2539 Win32Enforce!CloseHandle(handle); 2540 handle = null; 2541 } else version(Posix) { 2542 import unix = core.sys.posix.unistd; 2543 import core.sys.posix.fcntl; 2544 2545 ErrnoEnforce!(unix.close)(handle); 2546 handle = -1; 2547 } 2548 } 2549 } 2550 2551 /++ 2552 2553 +/ 2554 version(HasFile) class File : AbstractFile { 2555 2556 /++ 2557 Opens a file in synchronous access mode. 2558 2559 The permission mask is on used on posix systems FIXME: implement it 2560 +/ 2561 this(FilePath filename, OpenMode mode = OpenMode.readOnly, RequirePreexisting require = RequirePreexisting.no, uint specialFlags = 0, uint permMask = 0) { 2562 super(false, filename, mode, require, specialFlags); 2563 } 2564 2565 /++ 2566 2567 +/ 2568 ubyte[] read(scope ubyte[] buffer) { 2569 return null; 2570 } 2571 2572 /++ 2573 2574 +/ 2575 void write(in void[] buffer) { 2576 } 2577 2578 enum Seek { 2579 current, 2580 fromBeginning, 2581 fromEnd 2582 } 2583 2584 // Seeking/telling/sizing is not permitted when appending and some files don't support it 2585 // also not permitted in async mode 2586 void seek(long where, Seek fromWhence) {} 2587 long tell() { return 0; } 2588 } 2589 2590 /++ 2591 Only one operation can be pending at any time in the current implementation. 2592 +/ 2593 version(HasFile) class AsyncFile : AbstractFile { 2594 /++ 2595 Opens a file in asynchronous access mode. 2596 +/ 2597 this(FilePath filename, OpenMode mode = OpenMode.readOnly, RequirePreexisting require = RequirePreexisting.no, uint specialFlags = 0, uint permissionMask = 0) { 2598 // FIXME: implement permissionMask 2599 super(true, filename, mode, require, specialFlags); 2600 } 2601 2602 package(arsd) this(NativeFileHandle adoptPreSetup) { 2603 super(adoptPreSetup); 2604 } 2605 2606 /// 2607 AsyncReadRequest read(ubyte[] buffer, long offset = 0) { 2608 return new AsyncReadRequest(this, buffer, offset); 2609 } 2610 2611 /// 2612 AsyncWriteRequest write(const(void)[] buffer, long offset = 0) { 2613 return new AsyncWriteRequest(this, cast(ubyte[]) buffer, offset); 2614 } 2615 2616 } 2617 2618 /++ 2619 Reads or writes a file in one call. It might internally yield, but is generally blocking if it returns values. The callback ones depend on the implementation. 2620 2621 Tip: prefer the callback ones. If settings where async is possible, it will do async, and if not, it will sync. 2622 2623 NOT IMPLEMENTED 2624 +/ 2625 void writeFile(string filename, const(void)[] contents) { 2626 2627 } 2628 2629 /// ditto 2630 string readTextFile(string filename, string fileEncoding = null) { 2631 return null; 2632 } 2633 2634 /// ditto 2635 const(ubyte[]) readBinaryFile(string filename) { 2636 return null; 2637 } 2638 2639 /+ 2640 private Class recycleObject(Class, Args...)(Class objectToRecycle, Args args) { 2641 if(objectToRecycle is null) 2642 return new Class(args); 2643 // destroy nulls out the vtable which is the first thing in the object 2644 // so if it hasn't already been destroyed, we'll do it here 2645 if((*cast(void**) objectToRecycle) !is null) { 2646 assert(typeid(objectToRecycle) is typeid(Class)); // to make sure we're actually recycling the right kind of object 2647 .destroy(objectToRecycle); 2648 } 2649 2650 // then go ahead and reinitialize it 2651 ubyte[] rawData = (cast(ubyte*) cast(void*) objectToRecycle)[0 .. __traits(classInstanceSize, Class)]; 2652 rawData[] = (cast(ubyte[]) typeid(Class).initializer)[]; 2653 2654 objectToRecycle.__ctor(args); 2655 2656 return objectToRecycle; 2657 } 2658 +/ 2659 2660 /+ 2661 /++ 2662 Preallocates a class object without initializing it. 2663 2664 This is suitable *only* for passing to one of the functions in here that takes a preallocated object for recycling. 2665 +/ 2666 Class preallocate(Class)() { 2667 import core.memory; 2668 // FIXME: can i pass NO_SCAN here? 2669 return cast(Class) GC.calloc(__traits(classInstanceSize, Class), 0, typeid(Class)); 2670 } 2671 2672 OwnedClass!Class preallocateOnStack(Class)() { 2673 2674 } 2675 +/ 2676 2677 // thanks for a random person on stack overflow for this function 2678 version(Windows) 2679 BOOL MyCreatePipeEx( 2680 PHANDLE lpReadPipe, 2681 PHANDLE lpWritePipe, 2682 LPSECURITY_ATTRIBUTES lpPipeAttributes, 2683 DWORD nSize, 2684 DWORD dwReadMode, 2685 DWORD dwWriteMode 2686 ) 2687 { 2688 HANDLE ReadPipeHandle, WritePipeHandle; 2689 DWORD dwError; 2690 CHAR[MAX_PATH] PipeNameBuffer; 2691 2692 if (nSize == 0) { 2693 nSize = 4096; 2694 } 2695 2696 // FIXME: should be atomic op and gshared 2697 static shared(int) PipeSerialNumber = 0; 2698 2699 import core.stdc.string; 2700 import core.stdc.stdio; 2701 2702 sprintf(PipeNameBuffer.ptr, 2703 "\\\\.\\Pipe\\ArsdCoreAnonymousPipe.%08x.%08x".ptr, 2704 GetCurrentProcessId(), 2705 atomicOp!"+="(PipeSerialNumber, 1) 2706 ); 2707 2708 ReadPipeHandle = CreateNamedPipeA( 2709 PipeNameBuffer.ptr, 2710 1/*PIPE_ACCESS_INBOUND*/ | dwReadMode, 2711 0/*PIPE_TYPE_BYTE*/ | 0/*PIPE_WAIT*/, 2712 1, // Number of pipes 2713 nSize, // Out buffer size 2714 nSize, // In buffer size 2715 120 * 1000, // Timeout in ms 2716 lpPipeAttributes 2717 ); 2718 2719 if (! ReadPipeHandle) { 2720 return FALSE; 2721 } 2722 2723 WritePipeHandle = CreateFileA( 2724 PipeNameBuffer.ptr, 2725 GENERIC_WRITE, 2726 0, // No sharing 2727 lpPipeAttributes, 2728 OPEN_EXISTING, 2729 FILE_ATTRIBUTE_NORMAL | dwWriteMode, 2730 null // Template file 2731 ); 2732 2733 if (INVALID_HANDLE_VALUE == WritePipeHandle) { 2734 dwError = GetLastError(); 2735 CloseHandle( ReadPipeHandle ); 2736 SetLastError(dwError); 2737 return FALSE; 2738 } 2739 2740 *lpReadPipe = ReadPipeHandle; 2741 *lpWritePipe = WritePipeHandle; 2742 return( TRUE ); 2743 } 2744 2745 2746 2747 /+ 2748 2749 // this is probably useless. 2750 2751 /++ 2752 Creates a pair of anonymous pipes ready for async operations. 2753 2754 You can pass some preallocated objects to recycle if you like. 2755 +/ 2756 AsyncAnonymousPipe[2] anonymousPipePair(AsyncAnonymousPipe[2] preallocatedObjects = [null, null], bool inheritable = false) { 2757 version(Posix) { 2758 int[2] fds; 2759 auto ret = pipe(fds); 2760 2761 if(ret == -1) 2762 throw new SystemApiException("pipe", errno); 2763 2764 // FIXME: do we want them inheritable? and do we want both sides to be async? 2765 if(!inheritable) { 2766 setCloExec(fds[0]); 2767 setCloExec(fds[1]); 2768 } 2769 // if it is inherited, do we actually want it non-blocking? 2770 makeNonBlocking(fds[0]); 2771 makeNonBlocking(fds[1]); 2772 2773 return [ 2774 recycleObject(preallocatedObjects[0], fds[0]), 2775 recycleObject(preallocatedObjects[1], fds[1]), 2776 ]; 2777 } else version(Windows) { 2778 HANDLE rp, wp; 2779 // FIXME: do we want them inheritable? and do we want both sides to be async? 2780 if(!MyCreatePipeEx(&rp, &wp, null, 0, FILE_FLAG_OVERLAPPED, FILE_FLAG_OVERLAPPED)) 2781 throw new SystemApiException("MyCreatePipeEx", GetLastError()); 2782 return [ 2783 recycleObject(preallocatedObjects[0], rp), 2784 recycleObject(preallocatedObjects[1], wp), 2785 ]; 2786 } else throw ArsdException!"NotYetImplemented"(); 2787 } 2788 // on posix, just do pipe() w/ non block 2789 // on windows, do an overlapped named pipe server, connect, stop listening, return pair. 2790 +/ 2791 2792 /+ 2793 class NamedPipe : AsyncFile { 2794 2795 } 2796 +/ 2797 2798 /++ 2799 A named pipe ready to accept connections. 2800 2801 A Windows named pipe is an IPC mechanism usable on local machines or across a Windows network. 2802 +/ 2803 version(Windows) 2804 class NamedPipeServer { 2805 // unix domain socket or windows named pipe 2806 2807 // Promise!AsyncAnonymousPipe connect; 2808 // Promise!AsyncAnonymousPipe accept; 2809 2810 // when a new connection arrives, it calls your callback 2811 // can be on a specific thread or on any thread 2812 } 2813 2814 private version(Windows) extern(Windows) { 2815 const(char)* inet_ntop(int, const void*, char*, socklen_t); 2816 } 2817 2818 /++ 2819 Some functions that return arrays allow you to provide your own buffer. These are indicated in the type system as `UserProvidedBuffer!Type`, and you get to decide what you want to happen if the buffer is too small via the [OnOutOfSpace] parameter. 2820 2821 These are usually optional, since an empty user provided buffer with the default policy of reallocate will also work fine for whatever needs to be returned, thanks to the garbage collector taking care of it for you. 2822 2823 The API inside `UserProvidedBuffer` is all private to the arsd library implementation; your job is just to provide the buffer to it with [provideBuffer] or a constructor call and decide on your on-out-of-space policy. 2824 2825 $(TIP 2826 To properly size a buffer, I suggest looking at what covers about 80% of cases. Trying to cover everything often leads to wasted buffer space, and if you use a reallocate policy it can cover the rest. You might be surprised how far just two elements can go! 2827 ) 2828 2829 History: 2830 Added August 4, 2023 (dub v11.0) 2831 +/ 2832 struct UserProvidedBuffer(T) { 2833 private T[] buffer; 2834 private int actualLength; 2835 private OnOutOfSpace policy; 2836 2837 /++ 2838 2839 +/ 2840 public this(scope T[] buffer, OnOutOfSpace policy = OnOutOfSpace.reallocate) { 2841 this.buffer = buffer; 2842 this.policy = policy; 2843 } 2844 2845 package(arsd) bool append(T item) { 2846 if(actualLength < buffer.length) { 2847 buffer[actualLength++] = item; 2848 return true; 2849 } else final switch(policy) { 2850 case OnOutOfSpace.discard: 2851 return false; 2852 case OnOutOfSpace.exception: 2853 throw ArsdException!"Buffer out of space"(buffer.length, actualLength); 2854 case OnOutOfSpace.reallocate: 2855 buffer ~= item; 2856 actualLength++; 2857 return true; 2858 } 2859 } 2860 2861 package(arsd) T[] slice() return { 2862 return buffer[0 .. actualLength]; 2863 } 2864 } 2865 2866 /// ditto 2867 UserProvidedBuffer!T provideBuffer(T)(scope T[] buffer, OnOutOfSpace policy = OnOutOfSpace.reallocate) { 2868 return UserProvidedBuffer!T(buffer, policy); 2869 } 2870 2871 /++ 2872 Possible policies for [UserProvidedBuffer]s that run out of space. 2873 +/ 2874 enum OnOutOfSpace { 2875 reallocate, /// reallocate the buffer with the GC to make room 2876 discard, /// discard all contents that do not fit in your provided buffer 2877 exception, /// throw an exception if there is data that would not fit in your provided buffer 2878 } 2879 2880 2881 /++ 2882 For functions that give you an unknown address, you can use this to hold it. 2883 2884 Can get: 2885 ip4 2886 ip6 2887 unix 2888 abstract_ 2889 2890 name lookup for connect (stream or dgram) 2891 request canonical name? 2892 2893 interface lookup for bind (stream or dgram) 2894 +/ 2895 version(HasSocket) struct SocketAddress { 2896 import core.sys.posix.netdb; 2897 2898 /++ 2899 Provides the set of addresses to listen on all supported protocols on the machine for the given interfaces. `localhost` only listens on the loopback interface, whereas `allInterfaces` will listen on loopback as well as the others on the system (meaning it may be publicly exposed to the internet). 2900 2901 If you provide a buffer, I recommend using one of length two, so `SocketAddress[2]`, since this usually provides one address for ipv4 and one for ipv6. 2902 +/ 2903 static SocketAddress[] localhost(ushort port, return UserProvidedBuffer!SocketAddress buffer = null) { 2904 buffer.append(ip6("::1", port)); 2905 buffer.append(ip4("127.0.0.1", port)); 2906 return buffer.slice; 2907 } 2908 2909 /// ditto 2910 static SocketAddress[] allInterfaces(ushort port, return UserProvidedBuffer!SocketAddress buffer = null) { 2911 char[16] str; 2912 return allInterfaces(intToString(port, str[]), buffer); 2913 } 2914 2915 /// ditto 2916 static SocketAddress[] allInterfaces(scope const char[] serviceOrPort, return UserProvidedBuffer!SocketAddress buffer = null) { 2917 addrinfo hints; 2918 hints.ai_flags = AI_PASSIVE; 2919 hints.ai_socktype = SOCK_STREAM; // just to filter it down a little tbh 2920 return get(null, serviceOrPort, &hints, buffer); 2921 } 2922 2923 /++ 2924 Returns a single address object for the given protocol and parameters. 2925 2926 You probably should generally prefer [get], [localhost], or [allInterfaces] to have more flexible code. 2927 +/ 2928 static SocketAddress ip4(scope const char[] address, ushort port, bool forListening = false) { 2929 return getSingleAddress(AF_INET, AI_NUMERICHOST | (forListening ? AI_PASSIVE : 0), address, port); 2930 } 2931 2932 /// ditto 2933 static SocketAddress ip4(ushort port) { 2934 return ip4(null, port, true); 2935 } 2936 2937 /// ditto 2938 static SocketAddress ip6(scope const char[] address, ushort port, bool forListening = false) { 2939 return getSingleAddress(AF_INET6, AI_NUMERICHOST | (forListening ? AI_PASSIVE : 0), address, port); 2940 } 2941 2942 /// ditto 2943 static SocketAddress ip6(ushort port) { 2944 return ip6(null, port, true); 2945 } 2946 2947 /// ditto 2948 static SocketAddress unix(scope const char[] path) { 2949 // FIXME 2950 SocketAddress addr; 2951 return addr; 2952 } 2953 2954 /// ditto 2955 static SocketAddress abstract_(scope const char[] path) { 2956 char[190] buffer = void; 2957 buffer[0] = 0; 2958 buffer[1 .. path.length] = path[]; 2959 return unix(buffer[0 .. 1 + path.length]); 2960 } 2961 2962 private static SocketAddress getSingleAddress(int family, int flags, scope const char[] address, ushort port) { 2963 addrinfo hints; 2964 hints.ai_family = family; 2965 hints.ai_flags = flags; 2966 2967 char[16] portBuffer; 2968 char[] portString = intToString(port, portBuffer[]); 2969 2970 SocketAddress[1] addr; 2971 auto res = get(address, portString, &hints, provideBuffer(addr[])); 2972 if(res.length == 0) 2973 throw ArsdException!"bad address"(address.idup, port); 2974 return res[0]; 2975 } 2976 2977 /++ 2978 Calls `getaddrinfo` and returns the array of results. It will populate the data into the buffer you provide, if you provide one, otherwise it will allocate its own. 2979 +/ 2980 static SocketAddress[] get(scope const char[] nodeName, scope const char[] serviceOrPort, addrinfo* hints = null, return UserProvidedBuffer!SocketAddress buffer = null, scope bool delegate(scope addrinfo* ai) filter = null) @trusted { 2981 addrinfo* res; 2982 CharzBuffer node = nodeName; 2983 CharzBuffer service = serviceOrPort; 2984 auto ret = getaddrinfo(nodeName is null ? null : node.ptr, serviceOrPort is null ? null : service.ptr, hints, &res); 2985 if(ret == 0) { 2986 auto current = res; 2987 while(current) { 2988 if(filter is null || filter(current)) { 2989 SocketAddress addr; 2990 addr.addrlen = cast(socklen_t) current.ai_addrlen; 2991 switch(current.ai_family) { 2992 case AF_INET: 2993 addr.in4 = * cast(sockaddr_in*) current.ai_addr; 2994 break; 2995 case AF_INET6: 2996 addr.in6 = * cast(sockaddr_in6*) current.ai_addr; 2997 break; 2998 case AF_UNIX: 2999 addr.unix_address = * cast(sockaddr_un*) current.ai_addr; 3000 break; 3001 default: 3002 // skip 3003 } 3004 3005 if(!buffer.append(addr)) 3006 break; 3007 } 3008 3009 current = current.ai_next; 3010 } 3011 3012 freeaddrinfo(res); 3013 } else { 3014 version(Windows) { 3015 throw new WindowsApiException("getaddrinfo", ret); 3016 } else { 3017 const char* error = gai_strerror(ret); 3018 } 3019 } 3020 3021 return buffer.slice; 3022 } 3023 3024 /++ 3025 Returns a string representation of the address that identifies it in a custom format. 3026 3027 $(LIST 3028 * Unix domain socket addresses are their path prefixed with "unix:", unless they are in the abstract namespace, in which case it is prefixed with "abstract:" and the zero is trimmed out. For example, "unix:/tmp/pipe". 3029 3030 * IPv4 addresses are written in dotted decimal followed by a colon and the port number. For example, "127.0.0.1:8080". 3031 3032 * IPv6 addresses are written in colon separated hex format, but enclosed in brackets, then followed by the colon and port number. For example, "[::1]:8080". 3033 ) 3034 +/ 3035 string toString() const @trusted { 3036 char[200] buffer; 3037 switch(address.sa_family) { 3038 case AF_INET: 3039 auto writable = stringz(inet_ntop(address.sa_family, &in4.sin_addr, buffer.ptr, buffer.length)); 3040 auto it = writable.borrow; 3041 buffer[it.length] = ':'; 3042 auto numbers = intToString(port, buffer[it.length + 1 .. $]); 3043 return buffer[0 .. it.length + 1 + numbers.length].idup; 3044 case AF_INET6: 3045 buffer[0] = '['; 3046 auto writable = stringz(inet_ntop(address.sa_family, &in6.sin6_addr, buffer.ptr + 1, buffer.length - 1)); 3047 auto it = writable.borrow; 3048 buffer[it.length + 1] = ']'; 3049 buffer[it.length + 2] = ':'; 3050 auto numbers = intToString(port, buffer[it.length + 3 .. $]); 3051 return buffer[0 .. it.length + 3 + numbers.length].idup; 3052 case AF_UNIX: 3053 // FIXME: it might be abstract in which case stringz is wrong!!!!! 3054 auto writable = stringz(cast(char*) unix_address.sun_path.ptr).borrow; 3055 if(writable.length == 0) 3056 return "unix:"; 3057 string prefix = writable[0] == 0 ? "abstract:" : "unix:"; 3058 buffer[0 .. prefix.length] = prefix[]; 3059 buffer[prefix.length .. prefix.length + writable.length] = writable[writable[0] == 0 ? 1 : 0 .. $]; 3060 return buffer.idup; 3061 case AF_UNSPEC: 3062 return "<unspecified address>"; 3063 default: 3064 return "<unsupported address>"; // FIXME 3065 } 3066 } 3067 3068 ushort port() const @trusted { 3069 switch(address.sa_family) { 3070 case AF_INET: 3071 return ntohs(in4.sin_port); 3072 case AF_INET6: 3073 return ntohs(in6.sin6_port); 3074 default: 3075 return 0; 3076 } 3077 } 3078 3079 /+ 3080 @safe unittest { 3081 SocketAddress[4] buffer; 3082 foreach(addr; SocketAddress.get("arsdnet.net", "http", null, provideBuffer(buffer[]))) 3083 writeln(addr.toString()); 3084 } 3085 +/ 3086 3087 /+ 3088 unittest { 3089 // writeln(SocketAddress.ip4(null, 4444, true)); 3090 // writeln(SocketAddress.ip4("400.3.2.1", 4444)); 3091 // writeln(SocketAddress.ip4("bar", 4444)); 3092 foreach(addr; localhost(4444)) 3093 writeln(addr.toString()); 3094 } 3095 +/ 3096 3097 socklen_t addrlen = typeof(this).sizeof - socklen_t.sizeof; // the size of the union below 3098 3099 union { 3100 sockaddr address; 3101 3102 sockaddr_storage storage; 3103 3104 sockaddr_in in4; 3105 sockaddr_in6 in6; 3106 3107 sockaddr_un unix_address; 3108 } 3109 3110 /+ 3111 this(string node, string serviceOrPort, int family = 0) { 3112 // need to populate the approrpiate address and the length and make sure you set sa_family 3113 } 3114 +/ 3115 3116 int domain() { 3117 return address.sa_family; 3118 } 3119 sockaddr* rawAddr() return { 3120 return &address; 3121 } 3122 socklen_t rawAddrLength() { 3123 return addrlen; 3124 } 3125 3126 // FIXME it is AF_BLUETOOTH 3127 // see: https://people.csail.mit.edu/albert/bluez-intro/x79.html 3128 // see: https://learn.microsoft.com/en-us/windows/win32/Bluetooth/bluetooth-programming-with-windows-sockets 3129 } 3130 3131 private version(Windows) { 3132 struct sockaddr_un { 3133 ushort sun_family; 3134 char[108] sun_path; 3135 } 3136 } 3137 3138 version(HasFile) class AsyncSocket : AsyncFile { 3139 // otherwise: accept, bind, connect, shutdown, close. 3140 3141 static auto lastError() { 3142 version(Windows) 3143 return WSAGetLastError(); 3144 else 3145 return errno; 3146 } 3147 3148 static bool wouldHaveBlocked() { 3149 auto error = lastError; 3150 version(Windows) { 3151 return error == WSAEWOULDBLOCK || error == WSAETIMEDOUT; 3152 } else { 3153 return error == EAGAIN || error == EWOULDBLOCK; 3154 } 3155 } 3156 3157 version(Windows) 3158 enum INVALID = INVALID_SOCKET; 3159 else 3160 enum INVALID = -1; 3161 3162 // type is mostly SOCK_STREAM or SOCK_DGRAM 3163 /++ 3164 Creates a socket compatible with the given address. It does not actually connect or bind, nor store the address. You will want to pass it again to those functions: 3165 3166 --- 3167 auto socket = new Socket(address, Socket.Type.Stream); 3168 socket.connect(address).waitForCompletion(); 3169 --- 3170 +/ 3171 this(SocketAddress address, int type, int protocol = 0) { 3172 // need to look up these values for linux 3173 // type |= SOCK_NONBLOCK | SOCK_CLOEXEC; 3174 3175 handle_ = socket(address.domain(), type, protocol); 3176 if(handle == INVALID) 3177 throw new SystemApiException("socket", lastError()); 3178 3179 super(cast(NativeFileHandle) handle); // I think that cast is ok on Windows... i think 3180 3181 version(Posix) { 3182 makeNonBlocking(handle); 3183 setCloExec(handle); 3184 } 3185 3186 if(address.domain == AF_INET6) { 3187 int opt = 1; 3188 setsockopt(handle, IPPROTO_IPV6 /*SOL_IPV6*/, IPV6_V6ONLY, &opt, opt.sizeof); 3189 } 3190 3191 // FIXME: chekc for broadcast 3192 3193 // FIXME: REUSEADDR ? 3194 3195 // FIXME: also set NO_DELAY prolly 3196 // int opt = 1; 3197 // setsockopt(handle, IPPROTO_TCP, TCP_NODELAY, &opt, opt.sizeof); 3198 } 3199 3200 /++ 3201 Enabling NODELAY can give latency improvements if you are managing buffers on your end 3202 +/ 3203 void setNoDelay(bool enabled) { 3204 3205 } 3206 3207 /++ 3208 3209 `allowQuickRestart` will set the SO_REUSEADDR on unix and SO_DONTLINGER on Windows, 3210 allowing the application to be quickly restarted despite there still potentially being 3211 pending data in the tcp stack. 3212 3213 See https://stackoverflow.com/questions/3229860/what-is-the-meaning-of-so-reuseaddr-setsockopt-option-linux for more information. 3214 3215 If you already set your appropriate socket options or value correctness and reliability of the network stream over restart speed, leave this at the default `false`. 3216 +/ 3217 void bind(SocketAddress address, bool allowQuickRestart = false) { 3218 if(allowQuickRestart) { 3219 // FIXME 3220 } 3221 3222 auto ret = .bind(handle, address.rawAddr, address.rawAddrLength); 3223 if(ret == -1) 3224 throw new SystemApiException("bind", lastError); 3225 } 3226 3227 /++ 3228 You must call [bind] before this. 3229 3230 The backlog should be set to a value where your application can reliably catch up on the backlog in a reasonable amount of time under average load. It is meant to smooth over short duration bursts and making it too big will leave clients hanging - which might cause them to try to reconnect, thinking things got lost in transit, adding to your impossible backlog. 3231 3232 I personally tend to set this to be two per worker thread unless I have actual real world measurements saying to do something else. It is a bit arbitrary and not based on legitimate reasoning, it just seems to work for me (perhaps just because it has never really been put to the test). 3233 +/ 3234 void listen(int backlog) { 3235 auto ret = .listen(handle, backlog); 3236 if(ret == -1) 3237 throw new SystemApiException("listen", lastError); 3238 } 3239 3240 /++ 3241 +/ 3242 void shutdown(int how) { 3243 auto ret = .shutdown(handle, how); 3244 if(ret == -1) 3245 throw new SystemApiException("shutdown", lastError); 3246 } 3247 3248 /++ 3249 +/ 3250 override void close() { 3251 version(Windows) 3252 closesocket(handle); 3253 else 3254 .close(handle); 3255 handle_ = -1; 3256 } 3257 3258 /++ 3259 You can also construct your own request externally to control the memory more. 3260 +/ 3261 AsyncConnectRequest connect(SocketAddress address, ubyte[] bufferToSend = null) { 3262 return new AsyncConnectRequest(this, address, bufferToSend); 3263 } 3264 3265 /++ 3266 You can also construct your own request externally to control the memory more. 3267 +/ 3268 AsyncAcceptRequest accept() { 3269 return new AsyncAcceptRequest(this); 3270 } 3271 3272 // note that send is just sendto w/ a null address 3273 // and receive is just receivefrom w/ a null address 3274 /++ 3275 You can also construct your own request externally to control the memory more. 3276 +/ 3277 AsyncSendRequest send(const(ubyte)[] buffer, int flags = 0) { 3278 return new AsyncSendRequest(this, buffer, null, flags); 3279 } 3280 3281 /++ 3282 You can also construct your own request externally to control the memory more. 3283 +/ 3284 AsyncReceiveRequest receive(ubyte[] buffer, int flags = 0) { 3285 return new AsyncReceiveRequest(this, buffer, null, flags); 3286 } 3287 3288 /++ 3289 You can also construct your own request externally to control the memory more. 3290 +/ 3291 AsyncSendRequest sendTo(const(ubyte)[] buffer, SocketAddress* address, int flags = 0) { 3292 return new AsyncSendRequest(this, buffer, address, flags); 3293 } 3294 /++ 3295 You can also construct your own request externally to control the memory more. 3296 +/ 3297 AsyncReceiveRequest receiveFrom(ubyte[] buffer, SocketAddress* address, int flags = 0) { 3298 return new AsyncReceiveRequest(this, buffer, address, flags); 3299 } 3300 3301 /++ 3302 +/ 3303 SocketAddress localAddress() { 3304 SocketAddress addr; 3305 getsockname(handle, &addr.address, &addr.addrlen); 3306 return addr; 3307 } 3308 /++ 3309 +/ 3310 SocketAddress peerAddress() { 3311 SocketAddress addr; 3312 getpeername(handle, &addr.address, &addr.addrlen); 3313 return addr; 3314 } 3315 3316 // for unix sockets on unix only: send/receive fd, get peer creds 3317 3318 /++ 3319 3320 +/ 3321 final NativeSocketHandle handle() { 3322 return handle_; 3323 } 3324 3325 private NativeSocketHandle handle_; 3326 } 3327 3328 /++ 3329 Initiates a connection request and optionally sends initial data as soon as possible. 3330 3331 Calls `ConnectEx` on Windows and emulates it on other systems. 3332 3333 The entire buffer is sent before the operation is considered complete. 3334 3335 NOT IMPLEMENTED / NOT STABLE 3336 +/ 3337 version(HasSocket) class AsyncConnectRequest : AsyncOperationRequest { 3338 // FIXME: i should take a list of addresses and take the first one that succeeds, so a getaddrinfo can be sent straight in. 3339 this(AsyncSocket socket, SocketAddress address, ubyte[] dataToWrite) { 3340 3341 } 3342 3343 override void start() {} 3344 override void cancel() {} 3345 override bool isComplete() { return true; } 3346 override AsyncConnectResponse waitForCompletion() { assert(0); } 3347 } 3348 /++ 3349 +/ 3350 version(HasSocket) class AsyncConnectResponse : AsyncOperationResponse { 3351 const SystemErrorCode errorCode; 3352 3353 this(SystemErrorCode errorCode) { 3354 this.errorCode = errorCode; 3355 } 3356 3357 override bool wasSuccessful() { 3358 return errorCode.wasSuccessful; 3359 } 3360 3361 } 3362 3363 // FIXME: TransmitFile/sendfile support 3364 3365 /++ 3366 Calls `AcceptEx` on Windows and emulates it on other systems. 3367 3368 NOT IMPLEMENTED / NOT STABLE 3369 +/ 3370 version(HasSocket) class AsyncAcceptRequest : AsyncOperationRequest { 3371 AsyncSocket socket; 3372 3373 override void start() {} 3374 override void cancel() {} 3375 override bool isComplete() { return true; } 3376 override AsyncConnectResponse waitForCompletion() { assert(0); } 3377 3378 3379 struct LowLevelOperation { 3380 AsyncSocket file; 3381 ubyte[] buffer; 3382 SocketAddress* address; 3383 3384 this(typeof(this.tupleof) args) { 3385 this.tupleof = args; 3386 } 3387 3388 version(Windows) { 3389 auto opCall(OVERLAPPED* overlapped, LPOVERLAPPED_COMPLETION_ROUTINE ocr) { 3390 WSABUF buf; 3391 buf.len = cast(int) buffer.length; 3392 buf.buf = cast(typeof(buf.buf)) buffer.ptr; 3393 3394 uint flags; 3395 3396 if(address is null) 3397 return WSARecv(file.handle, &buf, 1, null, &flags, overlapped, ocr); 3398 else { 3399 return WSARecvFrom(file.handle, &buf, 1, null, &flags, &(address.address), &(address.addrlen), overlapped, ocr); 3400 } 3401 } 3402 } else { 3403 auto opCall() { 3404 int flags; 3405 if(address is null) 3406 return core.sys.posix.sys.socket.recv(file.handle, buffer.ptr, buffer.length, flags); 3407 else 3408 return core.sys.posix.sys.socket.recvfrom(file.handle, buffer.ptr, buffer.length, flags, &(address.address), &(address.addrlen)); 3409 } 3410 } 3411 3412 string errorString() { 3413 return "Receive"; 3414 } 3415 } 3416 mixin OverlappedIoRequest!(AsyncAcceptResponse, LowLevelOperation); 3417 3418 this(AsyncSocket socket, ubyte[] buffer = null, SocketAddress* address = null) { 3419 llo = LowLevelOperation(socket, buffer, address); 3420 this.response = typeof(this.response).defaultConstructed; 3421 } 3422 3423 // can also look up the local address 3424 } 3425 /++ 3426 +/ 3427 version(HasSocket) class AsyncAcceptResponse : AsyncOperationResponse { 3428 AsyncSocket newSocket; 3429 const SystemErrorCode errorCode; 3430 3431 this(SystemErrorCode errorCode, ubyte[] buffer) { 3432 this.errorCode = errorCode; 3433 } 3434 3435 this(AsyncSocket newSocket, SystemErrorCode errorCode) { 3436 this.newSocket = newSocket; 3437 this.errorCode = errorCode; 3438 } 3439 3440 override bool wasSuccessful() { 3441 return errorCode.wasSuccessful; 3442 } 3443 } 3444 3445 /++ 3446 +/ 3447 version(HasSocket) class AsyncReceiveRequest : AsyncOperationRequest { 3448 struct LowLevelOperation { 3449 AsyncSocket file; 3450 ubyte[] buffer; 3451 int flags; 3452 SocketAddress* address; 3453 3454 this(typeof(this.tupleof) args) { 3455 this.tupleof = args; 3456 } 3457 3458 version(Windows) { 3459 auto opCall(OVERLAPPED* overlapped, LPOVERLAPPED_COMPLETION_ROUTINE ocr) { 3460 WSABUF buf; 3461 buf.len = cast(int) buffer.length; 3462 buf.buf = cast(typeof(buf.buf)) buffer.ptr; 3463 3464 uint flags = this.flags; 3465 3466 if(address is null) 3467 return WSARecv(file.handle, &buf, 1, null, &flags, overlapped, ocr); 3468 else { 3469 return WSARecvFrom(file.handle, &buf, 1, null, &flags, &(address.address), &(address.addrlen), overlapped, ocr); 3470 } 3471 } 3472 } else { 3473 auto opCall() { 3474 if(address is null) 3475 return core.sys.posix.sys.socket.recv(file.handle, buffer.ptr, buffer.length, flags); 3476 else 3477 return core.sys.posix.sys.socket.recvfrom(file.handle, buffer.ptr, buffer.length, flags, &(address.address), &(address.addrlen)); 3478 } 3479 } 3480 3481 string errorString() { 3482 return "Receive"; 3483 } 3484 } 3485 mixin OverlappedIoRequest!(AsyncReceiveResponse, LowLevelOperation); 3486 3487 this(AsyncSocket socket, ubyte[] buffer, SocketAddress* address, int flags) { 3488 llo = LowLevelOperation(socket, buffer, flags, address); 3489 this.response = typeof(this.response).defaultConstructed; 3490 } 3491 3492 } 3493 /++ 3494 +/ 3495 version(HasSocket) class AsyncReceiveResponse : AsyncOperationResponse { 3496 const ubyte[] bufferWritten; 3497 const SystemErrorCode errorCode; 3498 3499 this(SystemErrorCode errorCode, const(ubyte)[] bufferWritten) { 3500 this.errorCode = errorCode; 3501 this.bufferWritten = bufferWritten; 3502 } 3503 3504 override bool wasSuccessful() { 3505 return errorCode.wasSuccessful; 3506 } 3507 } 3508 3509 /++ 3510 +/ 3511 version(HasSocket) class AsyncSendRequest : AsyncOperationRequest { 3512 struct LowLevelOperation { 3513 AsyncSocket file; 3514 const(ubyte)[] buffer; 3515 int flags; 3516 SocketAddress* address; 3517 3518 this(typeof(this.tupleof) args) { 3519 this.tupleof = args; 3520 } 3521 3522 version(Windows) { 3523 auto opCall(OVERLAPPED* overlapped, LPOVERLAPPED_COMPLETION_ROUTINE ocr) { 3524 WSABUF buf; 3525 buf.len = cast(int) buffer.length; 3526 buf.buf = cast(typeof(buf.buf)) buffer.ptr; 3527 3528 if(address is null) 3529 return WSASend(file.handle, &buf, 1, null, flags, overlapped, ocr); 3530 else { 3531 return WSASendTo(file.handle, &buf, 1, null, flags, address.rawAddr, address.rawAddrLength, overlapped, ocr); 3532 } 3533 } 3534 } else { 3535 auto opCall() { 3536 if(address is null) 3537 return core.sys.posix.sys.socket.send(file.handle, buffer.ptr, buffer.length, flags); 3538 else 3539 return core.sys.posix.sys.socket.sendto(file.handle, buffer.ptr, buffer.length, flags, address.rawAddr, address.rawAddrLength); 3540 } 3541 } 3542 3543 string errorString() { 3544 return "Send"; 3545 } 3546 } 3547 mixin OverlappedIoRequest!(AsyncSendResponse, LowLevelOperation); 3548 3549 this(AsyncSocket socket, const(ubyte)[] buffer, SocketAddress* address, int flags) { 3550 llo = LowLevelOperation(socket, buffer, flags, address); 3551 this.response = typeof(this.response).defaultConstructed; 3552 } 3553 } 3554 3555 /++ 3556 +/ 3557 version(HasSocket) class AsyncSendResponse : AsyncOperationResponse { 3558 const ubyte[] bufferWritten; 3559 const SystemErrorCode errorCode; 3560 3561 this(SystemErrorCode errorCode, const(ubyte)[] bufferWritten) { 3562 this.errorCode = errorCode; 3563 this.bufferWritten = bufferWritten; 3564 } 3565 3566 override bool wasSuccessful() { 3567 return errorCode.wasSuccessful; 3568 } 3569 3570 } 3571 3572 /++ 3573 A set of sockets bound and ready to accept connections on worker threads. 3574 3575 Depending on the specified address, it can be tcp, tcpv6, unix domain, or all of the above. 3576 3577 NOT IMPLEMENTED / NOT STABLE 3578 +/ 3579 version(HasSocket) class StreamServer { 3580 AsyncSocket[] sockets; 3581 3582 this(SocketAddress[] listenTo, int backlog = 8) { 3583 foreach(listen; listenTo) { 3584 auto socket = new AsyncSocket(listen, SOCK_STREAM); 3585 3586 // FIXME: allInterfaces for ipv6 also covers ipv4 so the bind can fail... 3587 // so we have to permit it to fail w/ address in use if we know we already 3588 // are listening to ipv6 3589 3590 // or there is a setsockopt ipv6 only thing i could set. 3591 3592 socket.bind(listen); 3593 socket.listen(backlog); 3594 sockets ~= socket; 3595 3596 // writeln(socket.localAddress.port); 3597 } 3598 3599 // i have to start accepting on each thread for each socket... 3600 } 3601 // when a new connection arrives, it calls your callback 3602 // can be on a specific thread or on any thread 3603 3604 3605 void start() { 3606 foreach(socket; sockets) { 3607 auto request = socket.accept(); 3608 request.start(); 3609 } 3610 } 3611 } 3612 3613 /+ 3614 unittest { 3615 auto ss = new StreamServer(SocketAddress.localhost(0)); 3616 } 3617 +/ 3618 3619 /++ 3620 A socket bound and ready to use receiveFrom 3621 3622 Depending on the address, it can be udp or unix domain. 3623 3624 NOT IMPLEMENTED / NOT STABLE 3625 +/ 3626 version(HasSocket) class DatagramListener { 3627 // whenever a udp message arrives, it calls your callback 3628 // can be on a specific thread or on any thread 3629 3630 // UDP is realistically just an async read on the bound socket 3631 // just it can get the "from" data out and might need the "more in packet" flag 3632 } 3633 3634 /++ 3635 Just in case I decide to change the implementation some day. 3636 +/ 3637 version(HasFile) alias AsyncAnonymousPipe = AsyncFile; 3638 3639 3640 // AsyncAnonymousPipe connectNamedPipe(AsyncAnonymousPipe preallocated, string name) 3641 3642 // unix fifos are considered just non-seekable files and have no special support in the lib; open them as a regular file w/ the async flag. 3643 3644 // DIRECTORY LISTINGS 3645 // not async, so if you want that, do it in a helper thread 3646 // just a convenient function to have (tho phobos has a decent one too, importing it expensive af) 3647 3648 /++ 3649 Note that the order of items called for your delegate is undefined; if you want it sorted, you'll have to collect and sort yourself. But it *might* be sorted by the OS (on Windows, it almost always is), so consider that when choosing a sorting algorithm. 3650 3651 History: 3652 previously in minigui as a private function. Moved to arsd.core on April 3, 2023 3653 +/ 3654 version(HasFile) GetFilesResult getFiles(string directory, scope void delegate(string name, bool isDirectory) dg) { 3655 // FIXME: my buffers here aren't great lol 3656 3657 SavedArgument[1] argsForException() { 3658 return [ 3659 SavedArgument("directory", LimitedVariant(directory)), 3660 ]; 3661 } 3662 3663 version(Windows) { 3664 WIN32_FIND_DATA data; 3665 // FIXME: if directory ends with / or \\ ? 3666 WCharzBuffer search = WCharzBuffer(directory ~ "/*"); 3667 auto handle = FindFirstFileW(search.ptr, &data); 3668 scope(exit) if(handle !is INVALID_HANDLE_VALUE) FindClose(handle); 3669 if(handle is INVALID_HANDLE_VALUE) { 3670 if(GetLastError() == ERROR_FILE_NOT_FOUND) 3671 return GetFilesResult.fileNotFound; 3672 throw new WindowsApiException("FindFirstFileW", GetLastError(), argsForException()[]); 3673 } 3674 3675 try_more: 3676 3677 string name = makeUtf8StringFromWindowsString(data.cFileName[0 .. findIndexOfZero(data.cFileName[])]); 3678 3679 dg(name, (data.dwFileAttributes & FILE_ATTRIBUTE_DIRECTORY) ? true : false); 3680 3681 auto ret = FindNextFileW(handle, &data); 3682 if(ret == 0) { 3683 if(GetLastError() == ERROR_NO_MORE_FILES) 3684 return GetFilesResult.success; 3685 throw new WindowsApiException("FindNextFileW", GetLastError(), argsForException()[]); 3686 } 3687 3688 goto try_more; 3689 3690 } else version(Posix) { 3691 import core.sys.posix.dirent; 3692 import core.stdc.errno; 3693 auto dir = opendir((directory ~ "\0").ptr); 3694 scope(exit) 3695 if(dir) closedir(dir); 3696 if(dir is null) 3697 throw new ErrnoApiException("opendir", errno, argsForException()); 3698 3699 auto dirent = readdir(dir); 3700 if(dirent is null) 3701 return GetFilesResult.fileNotFound; 3702 3703 try_more: 3704 3705 string name = dirent.d_name[0 .. findIndexOfZero(dirent.d_name[])].idup; 3706 3707 dg(name, dirent.d_type == DT_DIR); 3708 3709 dirent = readdir(dir); 3710 if(dirent is null) 3711 return GetFilesResult.success; 3712 3713 goto try_more; 3714 } else static assert(0); 3715 } 3716 3717 /// ditto 3718 enum GetFilesResult { 3719 success, 3720 fileNotFound 3721 } 3722 3723 /++ 3724 This is currently a simplified glob where only the * wildcard in the first or last position gets special treatment or a single * in the middle. 3725 3726 More things may be added later to be more like what Phobos supports. 3727 +/ 3728 bool matchesFilePattern(scope const(char)[] name, scope const(char)[] pattern) { 3729 if(pattern.length == 0) 3730 return false; 3731 if(pattern == "*") 3732 return true; 3733 if(pattern.length > 2 && pattern[0] == '*' && pattern[$-1] == '*') { 3734 // if the rest of pattern appears in name, it is good 3735 return name.indexOf(pattern[1 .. $-1]) != -1; 3736 } else if(pattern[0] == '*') { 3737 // if the rest of pattern is at end of name, it is good 3738 return name.endsWith(pattern[1 .. $]); 3739 } else if(pattern[$-1] == '*') { 3740 // if the rest of pattern is at start of name, it is good 3741 return name.startsWith(pattern[0 .. $-1]); 3742 } else if(pattern.length >= 3) { 3743 auto idx = pattern.indexOf("*"); 3744 if(idx != -1) { 3745 auto lhs = pattern[0 .. idx]; 3746 auto rhs = pattern[idx + 1 .. $]; 3747 if(name.length >= lhs.length + rhs.length) { 3748 return name.startsWith(lhs) && name.endsWith(rhs); 3749 } else { 3750 return false; 3751 } 3752 } 3753 } 3754 3755 return name == pattern; 3756 } 3757 3758 unittest { 3759 assert("test.html".matchesFilePattern("*")); 3760 assert("test.html".matchesFilePattern("*.html")); 3761 assert("test.html".matchesFilePattern("*.*")); 3762 assert("test.html".matchesFilePattern("test.*")); 3763 assert(!"test.html".matchesFilePattern("pest.*")); 3764 assert(!"test.html".matchesFilePattern("*.dhtml")); 3765 3766 assert("test.html".matchesFilePattern("t*.html")); 3767 assert(!"test.html".matchesFilePattern("e*.html")); 3768 } 3769 3770 package(arsd) int indexOf(scope const(char)[] haystack, scope const(char)[] needle) { 3771 if(haystack.length < needle.length) 3772 return -1; 3773 if(haystack == needle) 3774 return 0; 3775 foreach(i; 0 .. haystack.length - needle.length + 1) 3776 if(haystack[i .. i + needle.length] == needle) 3777 return cast(int) i; 3778 return -1; 3779 } 3780 3781 unittest { 3782 assert("foo".indexOf("f") == 0); 3783 assert("foo".indexOf("o") == 1); 3784 assert("foo".indexOf("foo") == 0); 3785 assert("foo".indexOf("oo") == 1); 3786 assert("foo".indexOf("fo") == 0); 3787 assert("foo".indexOf("boo") == -1); 3788 assert("foo".indexOf("food") == -1); 3789 } 3790 3791 package(arsd) bool endsWith(scope const(char)[] haystack, scope const(char)[] needle) { 3792 if(needle.length > haystack.length) 3793 return false; 3794 return haystack[$ - needle.length .. $] == needle; 3795 } 3796 3797 unittest { 3798 assert("foo".endsWith("o")); 3799 assert("foo".endsWith("oo")); 3800 assert("foo".endsWith("foo")); 3801 assert(!"foo".endsWith("food")); 3802 assert(!"foo".endsWith("d")); 3803 } 3804 3805 package(arsd) bool startsWith(scope const(char)[] haystack, scope const(char)[] needle) { 3806 if(needle.length > haystack.length) 3807 return false; 3808 return haystack[0 .. needle.length] == needle; 3809 } 3810 3811 unittest { 3812 assert("foo".startsWith("f")); 3813 assert("foo".startsWith("fo")); 3814 assert("foo".startsWith("foo")); 3815 assert(!"foo".startsWith("food")); 3816 assert(!"foo".startsWith("d")); 3817 } 3818 3819 3820 // FILE/DIR WATCHES 3821 // linux does it by name, windows and bsd do it by handle/descriptor 3822 // dispatches change event to either your thread or maybe the any task` queue. 3823 3824 /++ 3825 PARTIALLY IMPLEMENTED / NOT STABLE 3826 3827 +/ 3828 class DirectoryWatcher { 3829 private { 3830 version(Arsd_core_windows) { 3831 OVERLAPPED overlapped; 3832 HANDLE hDirectory; 3833 ubyte[] buffer; 3834 3835 extern(Windows) 3836 static void overlappedCompletionRoutine(DWORD dwErrorCode, DWORD dwNumberOfBytesTransferred, LPOVERLAPPED lpOverlapped) { 3837 typeof(this) rr = cast(typeof(this)) (cast(void*) lpOverlapped - typeof(this).overlapped.offsetof); 3838 3839 // dwErrorCode 3840 auto response = rr.buffer[0 .. dwNumberOfBytesTransferred]; 3841 3842 while(response.length) { 3843 auto fni = cast(FILE_NOTIFY_INFORMATION*) response.ptr; 3844 auto filename = fni.FileName[0 .. fni.FileNameLength]; 3845 3846 if(fni.NextEntryOffset) 3847 response = response[fni.NextEntryOffset .. $]; 3848 else 3849 response = response[$..$]; 3850 3851 // FIXME: I think I need to pin every overlapped op while it is pending 3852 // and unpin it when it is returned. GC.addRoot... but i don't wanna do that 3853 // every op so i guess i should do a refcount scheme similar to the other callback helper. 3854 3855 rr.changeHandler( 3856 FilePath(makeUtf8StringFromWindowsString(filename)), // FIXME: this is a relative path 3857 ChangeOperation.unknown // FIXME this is fni.Action 3858 ); 3859 } 3860 3861 rr.requestRead(); 3862 } 3863 3864 void requestRead() { 3865 DWORD ignored; 3866 if(!ReadDirectoryChangesW( 3867 hDirectory, 3868 buffer.ptr, 3869 cast(int) buffer.length, 3870 recursive, 3871 FILE_NOTIFY_CHANGE_LAST_WRITE | FILE_NOTIFY_CHANGE_CREATION | FILE_NOTIFY_CHANGE_FILE_NAME, 3872 &ignored, 3873 &overlapped, 3874 &overlappedCompletionRoutine 3875 )) { 3876 auto error = GetLastError(); 3877 /+ 3878 if(error == ERROR_IO_PENDING) { 3879 // not expected here, the docs say it returns true when queued 3880 } 3881 +/ 3882 3883 throw new SystemApiException("ReadDirectoryChangesW", error); 3884 } 3885 } 3886 } else version(Arsd_core_epoll) { 3887 static int inotifyfd = -1; // this is TLS since it is associated with the thread's event loop 3888 static ICoreEventLoop.UnregisterToken inotifyToken; 3889 static CallbackHelper inotifycb; 3890 static DirectoryWatcher[int] watchMappings; 3891 3892 static ~this() { 3893 if(inotifyfd != -1) { 3894 close(inotifyfd); 3895 inotifyfd = -1; 3896 } 3897 } 3898 3899 import core.sys.linux.sys.inotify; 3900 3901 int watchId = -1; 3902 3903 static void inotifyReady() { 3904 // read from it 3905 ubyte[256 /* NAME_MAX + 1 */ + inotify_event.sizeof] sbuffer; 3906 3907 auto ret = read(inotifyfd, sbuffer.ptr, sbuffer.length); 3908 if(ret == -1) { 3909 auto errno = errno; 3910 if(errno == EAGAIN || errno == EWOULDBLOCK) 3911 return; 3912 throw new SystemApiException("read inotify", errno); 3913 } else if(ret == 0) { 3914 assert(0, "I don't think this is ever supposed to happen"); 3915 } 3916 3917 auto buffer = sbuffer[0 .. ret]; 3918 3919 while(buffer.length > 0) { 3920 inotify_event* event = cast(inotify_event*) buffer.ptr; 3921 buffer = buffer[inotify_event.sizeof .. $]; 3922 char[] filename = cast(char[]) buffer[0 .. event.len]; 3923 buffer = buffer[event.len .. $]; 3924 3925 // note that filename is padded with zeroes, so it is actually a stringz 3926 3927 if(auto obj = event.wd in watchMappings) { 3928 (*obj).changeHandler( 3929 FilePath(stringz(filename.ptr).borrow.idup), // FIXME: this is a relative path 3930 ChangeOperation.unknown // FIXME 3931 ); 3932 } else { 3933 // it has probably already been removed 3934 } 3935 } 3936 } 3937 } else version(Arsd_core_kqueue) { 3938 int fd; 3939 CallbackHelper cb; 3940 } 3941 3942 FilePath path; 3943 string globPattern; 3944 bool recursive; 3945 void delegate(FilePath filename, ChangeOperation op) changeHandler; 3946 } 3947 3948 enum ChangeOperation { 3949 unknown, 3950 deleted, // NOTE_DELETE, IN_DELETE, FILE_NOTIFY_CHANGE_FILE_NAME 3951 written, // NOTE_WRITE / NOTE_EXTEND / NOTE_TRUNCATE, IN_MODIFY, FILE_NOTIFY_CHANGE_LAST_WRITE / FILE_NOTIFY_CHANGE_SIZE 3952 renamed, // NOTE_RENAME, the moved from/to in linux, FILE_NOTIFY_CHANGE_FILE_NAME 3953 metadataChanged // NOTE_ATTRIB, IN_ATTRIB, FILE_NOTIFY_CHANGE_ATTRIBUTES 3954 3955 // there is a NOTE_OPEN on freebsd 13, and the access change on Windows. and an open thing on linux. so maybe i can do note open/note_read too. 3956 } 3957 3958 /+ 3959 Windows and Linux work best when you watch directories. The operating system tells you the name of files as they change. 3960 3961 BSD doesn't support this. You can only get names and reports when a file is modified by watching specific files. AS such, when you watch a directory on those systems, your delegate will be called with a null path. Cross-platform applications should check for this and not assume the name is always usable. 3962 3963 inotify is kinda clearly the best of the bunch, with Windows in second place, and kqueue dead last. 3964 3965 3966 If path to watch is a directory, it signals when a file inside the directory (only one layer deep) is created or modified. This is the most efficient on Windows and Linux. 3967 3968 If a path is a file, it only signals when that specific file is written. This is most efficient on BSD. 3969 3970 3971 The delegate is called when something happens. Note that the path modified may not be accurate on all systems when you are watching a directory. 3972 +/ 3973 3974 /++ 3975 Watches a directory and its contents. If the `globPattern` is `null`, it will not attempt to add child items but also will not filter it, meaning you will be left with platform-specific behavior. 3976 3977 On Windows, the globPattern is just used to filter events. 3978 3979 On Linux, the `recursive` flag, if set, will cause it to add additional OS-level watches for each subdirectory. 3980 3981 On BSD, anything other than a null pattern will cause a directory scan to add files to the watch list. 3982 3983 For best results, use the most limited thing you need, as watches can get quite involved on the bsd systems. 3984 3985 Newly added files and subdirectories may not be automatically added in all cases, meaning if it is added and then subsequently modified, you might miss a notification. 3986 3987 If the event queue is too busy, the OS may skip a notification. 3988 3989 You should always offer some way for the user to force a refresh and not rely on notifications being present; they are a convenience when they work, not an always reliable method. 3990 +/ 3991 this(FilePath directoryToWatch, string globPattern, bool recursive, void delegate(FilePath pathModified, ChangeOperation op) dg) { 3992 this.path = directoryToWatch; 3993 this.globPattern = globPattern; 3994 this.recursive = recursive; 3995 this.changeHandler = dg; 3996 3997 version(Arsd_core_windows) { 3998 WCharzBuffer wname = directoryToWatch.path; 3999 buffer = new ubyte[](1024); 4000 hDirectory = CreateFileW( 4001 wname.ptr, 4002 GENERIC_READ, 4003 FILE_SHARE_READ, 4004 null, 4005 OPEN_EXISTING, 4006 FILE_ATTRIBUTE_NORMAL | FILE_FLAG_OVERLAPPED | FILE_FLAG_BACKUP_SEMANTICS, 4007 null 4008 ); 4009 if(hDirectory == INVALID_HANDLE_VALUE) 4010 throw new SystemApiException("CreateFileW", GetLastError()); 4011 4012 requestRead(); 4013 } else version(Arsd_core_epoll) { 4014 auto el = getThisThreadEventLoop(); 4015 4016 // no need for sync because it is thread-local 4017 if(inotifyfd == -1) { 4018 inotifyfd = inotify_init1(IN_NONBLOCK | IN_CLOEXEC); 4019 if(inotifyfd == -1) 4020 throw new SystemApiException("inotify_init1", errno); 4021 4022 inotifycb = new CallbackHelper(&inotifyReady); 4023 inotifyToken = el.addCallbackOnFdReadable(inotifyfd, inotifycb); 4024 } 4025 4026 uint event_mask = IN_CREATE | IN_MODIFY | IN_DELETE; // FIXME 4027 CharzBuffer dtw = directoryToWatch.path; 4028 auto watchId = inotify_add_watch(inotifyfd, dtw.ptr, event_mask); 4029 if(watchId < -1) 4030 throw new SystemApiException("inotify_add_watch", errno, [SavedArgument("path", LimitedVariant(directoryToWatch.path))]); 4031 4032 watchMappings[watchId] = this; 4033 4034 // FIXME: recursive needs to add child things individually 4035 4036 } else version(Arsd_core_kqueue) { 4037 auto el = cast(CoreEventLoopImplementation) getThisThreadEventLoop(); 4038 4039 // FIXME: need to scan for globPattern 4040 // when a new file is added, i'll have to diff my list to detect it and open it too 4041 // and recursive might need to scan down too. 4042 4043 kevent_t ev; 4044 4045 import core.sys.posix.fcntl; 4046 CharzBuffer buffer = CharzBuffer(directoryToWatch.path); 4047 fd = ErrnoEnforce!open(buffer.ptr, O_RDONLY); 4048 setCloExec(fd); 4049 4050 cb = new CallbackHelper(&triggered); 4051 4052 EV_SET(&ev, fd, EVFILT_VNODE, EV_ADD | EV_ENABLE | EV_CLEAR, NOTE_WRITE, 0, cast(void*) cb); 4053 ErrnoEnforce!kevent(el.kqueuefd, &ev, 1, null, 0, null); 4054 } else assert(0, "Not yet implemented for this platform"); 4055 } 4056 4057 private void triggered() { 4058 writeln("triggered"); 4059 } 4060 4061 void dispose() { 4062 version(Arsd_core_windows) { 4063 CloseHandle(hDirectory); 4064 } else version(Arsd_core_epoll) { 4065 watchMappings.remove(watchId); // I could also do this on the IN_IGNORE notification but idk 4066 inotify_rm_watch(inotifyfd, watchId); 4067 } else version(Arsd_core_kqueue) { 4068 ErrnoEnforce!close(fd); 4069 fd = -1; 4070 } 4071 } 4072 } 4073 4074 version(none) 4075 void main() { 4076 4077 // auto file = new AsyncFile(FilePath("test.txt"), AsyncFile.OpenMode.writeWithTruncation, AsyncFile.RequirePreexisting.yes); 4078 4079 /+ 4080 getFiles("c:/windows\\", (string filename, bool isDirectory) { 4081 writeln(filename, " ", isDirectory ? "[dir]": "[file]"); 4082 }); 4083 +/ 4084 4085 auto w = new DirectoryWatcher(FilePath("."), "*", false, (path, op) { 4086 writeln(path.path); 4087 }); 4088 getThisThreadEventLoop().run(() => false); 4089 } 4090 4091 /++ 4092 This starts up a local pipe. If it is already claimed, it just communicates with the existing one through the interface. 4093 +/ 4094 class SingleInstanceApplication { 4095 // FIXME 4096 } 4097 4098 version(none) 4099 void main() { 4100 4101 auto file = new AsyncFile(FilePath("test.txt"), AsyncFile.OpenMode.writeWithTruncation, AsyncFile.RequirePreexisting.yes); 4102 4103 auto buffer = cast(ubyte[]) "hello"; 4104 auto wr = new AsyncWriteRequest(file, buffer, 0); 4105 wr.start(); 4106 4107 wr.waitForCompletion(); 4108 4109 file.close(); 4110 } 4111 4112 /++ 4113 Implementation details of some requests. You shouldn't need to know any of this, the interface is all public. 4114 +/ 4115 mixin template OverlappedIoRequest(Response, LowLevelOperation) { 4116 private { 4117 LowLevelOperation llo; 4118 4119 OwnedClass!Response response; 4120 4121 version(Windows) { 4122 OVERLAPPED overlapped; 4123 4124 extern(Windows) 4125 static void overlappedCompletionRoutine(DWORD dwErrorCode, DWORD dwNumberOfBytesTransferred, LPOVERLAPPED lpOverlapped) { 4126 typeof(this) rr = cast(typeof(this)) (cast(void*) lpOverlapped - typeof(this).overlapped.offsetof); 4127 4128 rr.response = typeof(rr.response)(SystemErrorCode(dwErrorCode), rr.llo.buffer[0 .. dwNumberOfBytesTransferred]); 4129 rr.state_ = State.complete; 4130 4131 // FIXME: on complete? 4132 4133 // this will queue our CallbackHelper and that should be run at the end of the event loop after it is woken up by the APC run 4134 } 4135 } 4136 4137 version(Posix) { 4138 ICoreEventLoop.RearmToken eventRegistration; 4139 CallbackHelper cb; 4140 4141 final CallbackHelper getCb() { 4142 if(cb is null) 4143 cb = new CallbackHelper(&cbImpl); 4144 return cb; 4145 } 4146 4147 final void cbImpl() { 4148 // it is ready to complete, time to do it 4149 auto ret = llo(); 4150 markCompleted(ret, errno); 4151 } 4152 4153 void markCompleted(long ret, int errno) { 4154 // maybe i should queue an apc to actually do it, to ensure the event loop has cycled... FIXME 4155 if(ret == -1) 4156 response = typeof(response)(SystemErrorCode(errno), null); 4157 else 4158 response = typeof(response)(SystemErrorCode(0), llo.buffer[0 .. cast(size_t) ret]); 4159 state_ = State.complete; 4160 } 4161 } 4162 } 4163 4164 enum State { 4165 unused, 4166 started, 4167 inProgress, 4168 complete 4169 } 4170 private State state_; 4171 4172 override void start() { 4173 assert(state_ == State.unused); 4174 4175 state_ = State.started; 4176 4177 version(Windows) { 4178 if(llo(&overlapped, &overlappedCompletionRoutine)) { 4179 // all good, though GetLastError() might have some informative info 4180 } else { 4181 // operation failed, the operation is always ReadFileEx or WriteFileEx so it won't give the io pending thing here 4182 // should i issue error async? idk 4183 state_ = State.complete; 4184 throw new SystemApiException(llo.errorString(), GetLastError()); 4185 } 4186 4187 // ReadFileEx always queues, even if it completed synchronously. I *could* check the get overlapped result and sleepex here but i'm prolly better off just letting the event loop do its thing anyway. 4188 } else version(Posix) { 4189 4190 // first try to just do it 4191 auto ret = llo(); 4192 4193 auto errno = errno; 4194 if(ret == -1 && (errno == EAGAIN || errno == EWOULDBLOCK)) { // unable to complete right now, register and try when it is ready 4195 eventRegistration = getThisThreadEventLoop().addCallbackOnFdReadableOneShot(this.llo.file.handle, this.getCb); 4196 } else { 4197 // i could set errors sync or async and since it couldn't even start, i think a sync exception is the right way 4198 if(ret == -1) 4199 throw new SystemApiException(llo.errorString(), errno); 4200 markCompleted(ret, errno); // it completed synchronously (if it is an error nor not is handled by the completion handler) 4201 } 4202 } 4203 } 4204 4205 4206 override void cancel() { 4207 if(state_ == State.complete) 4208 return; // it has already finished, just leave it alone, no point discarding what is already done 4209 version(Windows) { 4210 if(state_ != State.unused) 4211 Win32Enforce!CancelIoEx(llo.file.AbstractFile.handle, &overlapped); 4212 // Windows will notify us when the cancellation is complete, so we need to wait for that before updating the state 4213 } else version(Posix) { 4214 if(state_ != State.unused) 4215 eventRegistration.unregister(); 4216 markCompleted(-1, ECANCELED); 4217 } 4218 } 4219 4220 override bool isComplete() { 4221 // just always let the event loop do it instead 4222 return state_ == State.complete; 4223 4224 /+ 4225 version(Windows) { 4226 return HasOverlappedIoCompleted(&overlapped); 4227 } else version(Posix) { 4228 return state_ == State.complete; 4229 4230 } 4231 +/ 4232 } 4233 4234 override Response waitForCompletion() { 4235 if(state_ == State.unused) 4236 start(); 4237 4238 // FIXME: if we are inside a fiber, we can set a oncomplete callback and then yield instead... 4239 if(state_ != State.complete) 4240 getThisThreadEventLoop().run(&isComplete); 4241 4242 /+ 4243 version(Windows) { 4244 SleepEx(INFINITE, true); 4245 4246 //DWORD numberTransferred; 4247 //Win32Enforce!GetOverlappedResult(file.handle, &overlapped, &numberTransferred, true); 4248 } else version(Posix) { 4249 getThisThreadEventLoop().run(&isComplete); 4250 } 4251 +/ 4252 4253 return response; 4254 } 4255 } 4256 4257 /++ 4258 You can write to a file asynchronously by creating one of these. 4259 +/ 4260 version(HasSocket) final class AsyncWriteRequest : AsyncOperationRequest { 4261 struct LowLevelOperation { 4262 AsyncFile file; 4263 ubyte[] buffer; 4264 long offset; 4265 4266 this(typeof(this.tupleof) args) { 4267 this.tupleof = args; 4268 } 4269 4270 version(Windows) { 4271 auto opCall(OVERLAPPED* overlapped, LPOVERLAPPED_COMPLETION_ROUTINE ocr) { 4272 overlapped.Offset = (cast(ulong) offset) & 0xffff_ffff; 4273 overlapped.OffsetHigh = ((cast(ulong) offset) >> 32) & 0xffff_ffff; 4274 return WriteFileEx(file.handle, buffer.ptr, cast(int) buffer.length, overlapped, ocr); 4275 } 4276 } else { 4277 auto opCall() { 4278 return core.sys.posix.unistd.write(file.handle, buffer.ptr, buffer.length); 4279 } 4280 } 4281 4282 string errorString() { 4283 return "Write"; 4284 } 4285 } 4286 mixin OverlappedIoRequest!(AsyncWriteResponse, LowLevelOperation); 4287 4288 this(AsyncFile file, ubyte[] buffer, long offset) { 4289 this.llo = LowLevelOperation(file, buffer, offset); 4290 response = typeof(response).defaultConstructed; 4291 } 4292 } 4293 4294 /++ 4295 4296 +/ 4297 class AsyncWriteResponse : AsyncOperationResponse { 4298 const ubyte[] bufferWritten; 4299 const SystemErrorCode errorCode; 4300 4301 this(SystemErrorCode errorCode, const(ubyte)[] bufferWritten) { 4302 this.errorCode = errorCode; 4303 this.bufferWritten = bufferWritten; 4304 } 4305 4306 override bool wasSuccessful() { 4307 return errorCode.wasSuccessful; 4308 } 4309 } 4310 4311 /++ 4312 4313 +/ 4314 version(HasSocket) final class AsyncReadRequest : AsyncOperationRequest { 4315 struct LowLevelOperation { 4316 AsyncFile file; 4317 ubyte[] buffer; 4318 long offset; 4319 4320 this(typeof(this.tupleof) args) { 4321 this.tupleof = args; 4322 } 4323 4324 version(Windows) { 4325 auto opCall(OVERLAPPED* overlapped, LPOVERLAPPED_COMPLETION_ROUTINE ocr) { 4326 overlapped.Offset = (cast(ulong) offset) & 0xffff_ffff; 4327 overlapped.OffsetHigh = ((cast(ulong) offset) >> 32) & 0xffff_ffff; 4328 return ReadFileEx(file.handle, buffer.ptr, cast(int) buffer.length, overlapped, ocr); 4329 } 4330 } else { 4331 auto opCall() { 4332 return core.sys.posix.unistd.read(file.handle, buffer.ptr, buffer.length); 4333 } 4334 } 4335 4336 string errorString() { 4337 return "Read"; 4338 } 4339 } 4340 mixin OverlappedIoRequest!(AsyncReadResponse, LowLevelOperation); 4341 4342 /++ 4343 The file must have the overlapped flag enabled on Windows and the nonblock flag set on Posix. 4344 4345 The buffer MUST NOT be touched by you - not used by another request, modified, read, or freed, including letting a static array going out of scope - until this request's `isComplete` returns `true`. 4346 4347 The offset is where to start reading a disk file. For all other types of files, pass 0. 4348 +/ 4349 this(AsyncFile file, ubyte[] buffer, long offset) { 4350 this.llo = LowLevelOperation(file, buffer, offset); 4351 response = typeof(response).defaultConstructed; 4352 } 4353 4354 /++ 4355 4356 +/ 4357 // abstract void repeat(); 4358 } 4359 4360 /++ 4361 4362 +/ 4363 class AsyncReadResponse : AsyncOperationResponse { 4364 const ubyte[] bufferRead; 4365 const SystemErrorCode errorCode; 4366 4367 this(SystemErrorCode errorCode, const(ubyte)[] bufferRead) { 4368 this.errorCode = errorCode; 4369 this.bufferRead = bufferRead; 4370 } 4371 4372 override bool wasSuccessful() { 4373 return errorCode.wasSuccessful; 4374 } 4375 } 4376 4377 /+ 4378 Tasks: 4379 startTask() 4380 startSubTask() - what if it just did this when it knows it is being run from inside a task? 4381 runHelperFunction() - whomever it reports to is the parent 4382 +/ 4383 4384 version(HasThread) class ScheduableTask : Fiber { 4385 private void delegate() dg; 4386 4387 // linked list stuff 4388 private static ScheduableTask taskRoot; 4389 private ScheduableTask previous; 4390 private ScheduableTask next; 4391 4392 // need the controlling thread to know how to wake it up if it receives a message 4393 private Thread controllingThread; 4394 4395 // the api 4396 4397 this(void delegate() dg) { 4398 assert(dg !is null); 4399 4400 this.dg = dg; 4401 super(&taskRunner); 4402 4403 if(taskRoot !is null) { 4404 this.next = taskRoot; 4405 taskRoot.previous = this; 4406 } 4407 taskRoot = this; 4408 } 4409 4410 /+ 4411 enum BehaviorOnCtrlC { 4412 ignore, 4413 cancel, 4414 deliverMessage 4415 } 4416 +/ 4417 4418 private bool cancelled; 4419 4420 public void cancel() { 4421 this.cancelled = true; 4422 // if this is running, we can throw immediately 4423 // otherwise if we're calling from an appropriate thread, we can call it immediately 4424 // otherwise we need to queue a wakeup to its own thread. 4425 // tbh we should prolly just queue it every time 4426 } 4427 4428 private void taskRunner() { 4429 try { 4430 dg(); 4431 } catch(TaskCancelledException tce) { 4432 // this space intentionally left blank; 4433 // the purpose of this exception is to just 4434 // let the fiber's destructors run before we 4435 // let it die. 4436 } catch(Throwable t) { 4437 if(taskUncaughtException is null) { 4438 throw t; 4439 } else { 4440 taskUncaughtException(t); 4441 } 4442 } finally { 4443 if(this is taskRoot) { 4444 taskRoot = taskRoot.next; 4445 if(taskRoot !is null) 4446 taskRoot.previous = null; 4447 } else { 4448 assert(this.previous !is null); 4449 assert(this.previous.next is this); 4450 this.previous.next = this.next; 4451 if(this.next !is null) 4452 this.next.previous = this.previous; 4453 } 4454 } 4455 } 4456 } 4457 4458 /++ 4459 4460 +/ 4461 void delegate(Throwable t) taskUncaughtException; 4462 4463 /++ 4464 Gets an object that lets you control a schedulable task (which is a specialization of a fiber) and can be used in an `if` statement. 4465 4466 --- 4467 if(auto controller = inSchedulableTask()) { 4468 controller.yieldUntilReadable(...); 4469 } 4470 --- 4471 4472 History: 4473 Added August 11, 2023 (dub v11.1) 4474 +/ 4475 version(HasThread) SchedulableTaskController inSchedulableTask() { 4476 import core.thread.fiber; 4477 4478 if(auto fiber = Fiber.getThis) { 4479 return SchedulableTaskController(cast(ScheduableTask) fiber); 4480 } 4481 4482 return SchedulableTaskController(null); 4483 } 4484 4485 /// ditto 4486 version(HasThread) struct SchedulableTaskController { 4487 private this(ScheduableTask fiber) { 4488 this.fiber = fiber; 4489 } 4490 4491 private ScheduableTask fiber; 4492 4493 /++ 4494 4495 +/ 4496 bool opCast(T : bool)() { 4497 return fiber !is null; 4498 } 4499 4500 /++ 4501 4502 +/ 4503 version(Posix) 4504 void yieldUntilReadable(NativeFileHandle handle) { 4505 assert(fiber !is null); 4506 4507 auto cb = new CallbackHelper(() { fiber.call(); }); 4508 4509 // FIXME: if the fd is already registered in this thread it can throw... 4510 version(Windows) 4511 auto rearmToken = getThisThreadEventLoop().addCallbackOnFdReadableOneShot(handle, cb); 4512 else 4513 auto rearmToken = getThisThreadEventLoop().addCallbackOnFdReadableOneShot(handle, cb); 4514 4515 // FIXME: this is only valid if the fiber is only ever going to run in this thread! 4516 fiber.yield(); 4517 4518 rearmToken.unregister(); 4519 4520 // what if there are other messages, like a ctrl+c? 4521 if(fiber.cancelled) 4522 throw new TaskCancelledException(); 4523 } 4524 4525 version(Windows) 4526 void yieldUntilSignaled(NativeFileHandle handle) { 4527 // add it to the WaitForMultipleObjects thing w/ a cb 4528 } 4529 } 4530 4531 class TaskCancelledException : object.Exception { 4532 this() { 4533 super("Task cancelled"); 4534 } 4535 } 4536 4537 version(HasThread) private class CoreWorkerThread : Thread { 4538 this(EventLoopType type) { 4539 this.type = type; 4540 4541 // task runners are supposed to have smallish stacks since they either just run a single callback or call into fibers 4542 // the helper runners might be a bit bigger tho 4543 super(&run); 4544 } 4545 void run() { 4546 eventLoop = getThisThreadEventLoop(this.type); 4547 atomicOp!"+="(startedCount, 1); 4548 atomicOp!"+="(runningCount, 1); 4549 scope(exit) { 4550 atomicOp!"-="(runningCount, 1); 4551 } 4552 4553 eventLoop.run(() => cancelled); 4554 } 4555 4556 private bool cancelled; 4557 4558 void cancel() { 4559 cancelled = true; 4560 } 4561 4562 EventLoopType type; 4563 ICoreEventLoop eventLoop; 4564 4565 __gshared static { 4566 CoreWorkerThread[] taskRunners; 4567 CoreWorkerThread[] helperRunners; 4568 ICoreEventLoop mainThreadLoop; 4569 4570 // for the helper function thing on the bsds i could have my own little circular buffer of availability 4571 4572 shared(int) startedCount; 4573 shared(int) runningCount; 4574 4575 bool started; 4576 4577 void setup(int numberOfTaskRunners, int numberOfHelpers) { 4578 assert(!started); 4579 synchronized { 4580 mainThreadLoop = getThisThreadEventLoop(); 4581 4582 foreach(i; 0 .. numberOfTaskRunners) { 4583 auto nt = new CoreWorkerThread(EventLoopType.TaskRunner); 4584 taskRunners ~= nt; 4585 nt.start(); 4586 } 4587 foreach(i; 0 .. numberOfHelpers) { 4588 auto nt = new CoreWorkerThread(EventLoopType.HelperWorker); 4589 helperRunners ~= nt; 4590 nt.start(); 4591 } 4592 4593 const expectedCount = numberOfHelpers + numberOfTaskRunners; 4594 4595 while(startedCount < expectedCount) { 4596 Thread.yield(); 4597 } 4598 4599 started = true; 4600 } 4601 } 4602 4603 void cancelAll() { 4604 foreach(runner; taskRunners) 4605 runner.cancel(); 4606 foreach(runner; helperRunners) 4607 runner.cancel(); 4608 4609 } 4610 } 4611 } 4612 4613 private int numberOfCpus() { 4614 return 4; // FIXME 4615 } 4616 4617 /++ 4618 To opt in to the full functionality of this module with customization opportunity, create one and only one of these objects that is valid for exactly the lifetime of the application. 4619 4620 Normally, this means writing a main like this: 4621 4622 --- 4623 import arsd.core; 4624 void main() { 4625 ArsdCoreApplication app = ArsdCoreApplication("Your app name"); 4626 4627 // do your setup here 4628 4629 // the rest of your code here 4630 } 4631 --- 4632 4633 Its destructor runs the event loop then waits to for the workers to finish to clean them up. 4634 +/ 4635 // FIXME: single instance? 4636 version(HasThread) struct ArsdCoreApplication { 4637 private ICoreEventLoop impl; 4638 4639 /++ 4640 default number of threads is to split your cpus between blocking function runners and task runners 4641 +/ 4642 this(string applicationName) { 4643 auto num = numberOfCpus(); 4644 num /= 2; 4645 if(num <= 0) 4646 num = 1; 4647 this(applicationName, num, num); 4648 } 4649 4650 /++ 4651 4652 +/ 4653 this(string applicationName, int numberOfTaskRunners, int numberOfHelpers) { 4654 impl = getThisThreadEventLoop(EventLoopType.Explicit); 4655 CoreWorkerThread.setup(numberOfTaskRunners, numberOfHelpers); 4656 } 4657 4658 @disable this(); 4659 @disable this(this); 4660 /++ 4661 This must be deterministically destroyed. 4662 +/ 4663 @disable new(); 4664 4665 ~this() { 4666 if(!alreadyRun) 4667 run(); 4668 exitApplication(); 4669 waitForWorkersToExit(3000); 4670 } 4671 4672 void exitApplication() { 4673 CoreWorkerThread.cancelAll(); 4674 } 4675 4676 void waitForWorkersToExit(int timeoutMilliseconds) { 4677 4678 } 4679 4680 private bool alreadyRun; 4681 4682 void run() { 4683 impl.run(() => false); 4684 alreadyRun = true; 4685 } 4686 } 4687 4688 4689 private class CoreEventLoopImplementation : ICoreEventLoop { 4690 version(EmptyEventLoop) void runOnce(){} 4691 version(EmptyCoreEvent) 4692 { 4693 UnregisterToken addCallbackOnFdReadable(int fd, CallbackHelper cb){return typeof(return).init;} 4694 RearmToken addCallbackOnFdReadableOneShot(int fd, CallbackHelper cb){return typeof(return).init;} 4695 RearmToken addCallbackOnFdWritableOneShot(int fd, CallbackHelper cb){return typeof(return).init;} 4696 private void rearmFd(RearmToken token) {} 4697 } 4698 4699 version(Arsd_core_kqueue) { 4700 // this thread apc dispatches go as a custom event to the queue 4701 // the other queues go through one byte at a time pipes (barf). freebsd 13 and newest nbsd have eventfd too tho so maybe i can use them but the other kqueue systems don't. 4702 4703 void runOnce() { 4704 kevent_t[16] ev; 4705 //timespec tout = timespec(1, 0); 4706 auto nev = kevent(kqueuefd, null, 0, ev.ptr, ev.length, null/*&tout*/); 4707 if(nev == -1) { 4708 // FIXME: EINTR 4709 throw new SystemApiException("kevent", errno); 4710 } else if(nev == 0) { 4711 // timeout 4712 } else { 4713 foreach(event; ev[0 .. nev]) { 4714 if(event.filter == EVFILT_SIGNAL) { 4715 // FIXME: I could prolly do this better tbh 4716 markSignalOccurred(cast(int) event.ident); 4717 signalChecker(); 4718 } else { 4719 // FIXME: event.filter more specific? 4720 CallbackHelper cb = cast(CallbackHelper) event.udata; 4721 cb.call(); 4722 } 4723 } 4724 } 4725 } 4726 4727 // FIXME: idk how to make one event that multiple kqueues can listen to w/o being shared 4728 // maybe a shared kqueue could work that the thread kqueue listen to (which i rejected for 4729 // epoll cuz it caused thundering herd problems but maybe it'd work here) 4730 4731 UnregisterToken addCallbackOnFdReadable(int fd, CallbackHelper cb) { 4732 kevent_t ev; 4733 4734 EV_SET(&ev, fd, EVFILT_READ, EV_ADD | EV_ENABLE/* | EV_ONESHOT*/, 0, 0, cast(void*) cb); 4735 4736 ErrnoEnforce!kevent(kqueuefd, &ev, 1, null, 0, null); 4737 4738 return UnregisterToken(this, fd, cb); 4739 } 4740 4741 RearmToken addCallbackOnFdReadableOneShot(int fd, CallbackHelper cb) { 4742 kevent_t ev; 4743 4744 EV_SET(&ev, fd, EVFILT_READ, EV_ADD | EV_ENABLE/* | EV_ONESHOT*/, 0, 0, cast(void*) cb); 4745 4746 ErrnoEnforce!kevent(kqueuefd, &ev, 1, null, 0, null); 4747 4748 return RearmToken(true, this, fd, cb, 0); 4749 } 4750 4751 RearmToken addCallbackOnFdWritableOneShot(int fd, CallbackHelper cb) { 4752 kevent_t ev; 4753 4754 EV_SET(&ev, fd, EVFILT_WRITE, EV_ADD | EV_ENABLE/* | EV_ONESHOT*/, 0, 0, cast(void*) cb); 4755 4756 ErrnoEnforce!kevent(kqueuefd, &ev, 1, null, 0, null); 4757 4758 return RearmToken(false, this, fd, cb, 0); 4759 } 4760 4761 private void rearmFd(RearmToken token) { 4762 if(token.readable) 4763 cast(void) addCallbackOnFdReadableOneShot(token.fd, token.cb); 4764 else 4765 cast(void) addCallbackOnFdWritableOneShot(token.fd, token.cb); 4766 } 4767 4768 private void triggerGlobalEvent() { 4769 ubyte a; 4770 import core.sys.posix.unistd; 4771 write(kqueueGlobalFd[1], &a, 1); 4772 } 4773 4774 private this() { 4775 kqueuefd = ErrnoEnforce!kqueue(); 4776 setCloExec(kqueuefd); // FIXME O_CLOEXEC 4777 4778 if(kqueueGlobalFd[0] == 0) { 4779 import core.sys.posix.unistd; 4780 pipe(kqueueGlobalFd); 4781 setCloExec(kqueueGlobalFd[0]); 4782 setCloExec(kqueueGlobalFd[1]); 4783 4784 signal(SIGINT, SIG_IGN); // FIXME 4785 } 4786 4787 kevent_t ev; 4788 4789 EV_SET(&ev, SIGCHLD, EVFILT_SIGNAL, EV_ADD | EV_ENABLE, 0, 0, null); 4790 ErrnoEnforce!kevent(kqueuefd, &ev, 1, null, 0, null); 4791 EV_SET(&ev, SIGINT, EVFILT_SIGNAL, EV_ADD | EV_ENABLE, 0, 0, null); 4792 ErrnoEnforce!kevent(kqueuefd, &ev, 1, null, 0, null); 4793 4794 globalEventSent = new CallbackHelper(&readGlobalEvent); 4795 EV_SET(&ev, kqueueGlobalFd[0], EVFILT_READ, EV_ADD | EV_ENABLE, 0, 0, cast(void*) globalEventSent); 4796 ErrnoEnforce!kevent(kqueuefd, &ev, 1, null, 0, null); 4797 } 4798 4799 private int kqueuefd = -1; 4800 4801 private CallbackHelper globalEventSent; 4802 void readGlobalEvent() { 4803 kevent_t event; 4804 4805 import core.sys.posix.unistd; 4806 ubyte a; 4807 read(kqueueGlobalFd[0], &a, 1); 4808 4809 // FIXME: the thread is woken up, now we need to check the circualr buffer queue 4810 } 4811 4812 private __gshared int[2] kqueueGlobalFd; 4813 } 4814 4815 /+ 4816 // this setup needs no extra allocation 4817 auto op = read(file, buffer); 4818 op.oncomplete = &thisfiber.call; 4819 op.start(); 4820 thisfiber.yield(); 4821 auto result = op.waitForCompletion(); // guaranteed to return instantly thanks to previous setup 4822 4823 can generically abstract that into: 4824 4825 auto result = thisTask.await(read(file, buffer)); 4826 4827 4828 You MUST NOT use buffer in any way - not read, modify, deallocate, reuse, anything - until the PendingOperation is complete. 4829 4830 Note that PendingOperation may just be a wrapper around an internally allocated object reference... but then if you do a waitForFirstToComplete what happens? 4831 4832 those could of course just take the value type things 4833 +/ 4834 4835 4836 version(Arsd_core_windows) { 4837 // all event loops share the one iocp, Windows 4838 // manages how to do it 4839 __gshared HANDLE iocpTaskRunners; 4840 __gshared HANDLE iocpWorkers; 4841 4842 HANDLE[] handles; 4843 4844 // i think to terminate i just have to post the message at least once for every thread i know about, maybe a few more times for threads i don't know about. 4845 4846 bool isWorker; // if it is a worker we wait on the iocp, if not we wait on msg 4847 4848 void runOnce() { 4849 if(isWorker) { 4850 // this function is only supported on Windows Vista and up, so using this 4851 // means dropping support for XP. 4852 //GetQueuedCompletionStatusEx(); 4853 assert(0); // FIXME 4854 } else { 4855 auto wto = 0; 4856 4857 auto waitResult = MsgWaitForMultipleObjectsEx( 4858 cast(int) handles.length, handles.ptr, 4859 (wto == 0 ? INFINITE : wto), /* timeout */ 4860 0x04FF, /* QS_ALLINPUT */ 4861 0x0002 /* MWMO_ALERTABLE */ | 0x0004 /* MWMO_INPUTAVAILABLE */); 4862 4863 enum WAIT_OBJECT_0 = 0; 4864 if(waitResult >= WAIT_OBJECT_0 && waitResult < handles.length + WAIT_OBJECT_0) { 4865 auto h = handles[waitResult - WAIT_OBJECT_0]; 4866 // FIXME: run the handle ready callback 4867 } else if(waitResult == handles.length + WAIT_OBJECT_0) { 4868 // message ready 4869 int count; 4870 MSG message; 4871 while(PeekMessage(&message, null, 0, 0, PM_NOREMOVE)) { // need to peek since sometimes MsgWaitForMultipleObjectsEx returns even though GetMessage can block. tbh i don't fully understand it but the docs say it is foreground activation 4872 auto ret = GetMessage(&message, null, 0, 0); 4873 if(ret == -1) 4874 throw new WindowsApiException("GetMessage", GetLastError()); 4875 TranslateMessage(&message); 4876 DispatchMessage(&message); 4877 4878 count++; 4879 if(count > 10) 4880 break; // take the opportunity to catch up on other events 4881 4882 if(ret == 0) { // WM_QUIT 4883 // EventLoop.quitApplication(); 4884 assert(0); // FIXME 4885 //break; 4886 } 4887 } 4888 } else if(waitResult == 0x000000C0L /* WAIT_IO_COMPLETION */) { 4889 SleepEx(0, true); // I call this to give it a chance to do stuff like async io 4890 } else if(waitResult == 258L /* WAIT_TIMEOUT */) { 4891 // timeout, should never happen since we aren't using it 4892 } else if(waitResult == 0xFFFFFFFF) { 4893 // failed 4894 throw new WindowsApiException("MsgWaitForMultipleObjectsEx", GetLastError()); 4895 } else { 4896 // idk.... 4897 } 4898 } 4899 } 4900 } 4901 4902 version(Posix) { 4903 private __gshared uint sigChildHappened = 0; 4904 private __gshared uint sigIntrHappened = 0; 4905 4906 static void signalChecker() { 4907 if(cas(&sigChildHappened, 1, 0)) { 4908 while(true) { // multiple children could have exited before we processed the notification 4909 4910 import core.sys.posix.sys.wait; 4911 4912 int status; 4913 auto pid = waitpid(-1, &status, WNOHANG); 4914 if(pid == -1) { 4915 import core.stdc.errno; 4916 auto errno = errno; 4917 if(errno == ECHILD) 4918 break; // also all done, there are no children left 4919 // no need to check EINTR since we set WNOHANG 4920 throw new ErrnoApiException("waitpid", errno); 4921 } 4922 if(pid == 0) 4923 break; // all done, all children are still running 4924 4925 // look up the pid for one of our objects 4926 // if it is found, inform it of its status 4927 // and then inform its controlling thread 4928 // to wake up so it can check its waitForCompletion, 4929 // trigger its callbacks, etc. 4930 4931 ExternalProcess.recordChildTerminated(pid, status); 4932 } 4933 4934 } 4935 if(cas(&sigIntrHappened, 1, 0)) { 4936 // FIXME 4937 import core.stdc.stdlib; 4938 exit(0); 4939 } 4940 } 4941 4942 /++ 4943 Informs the arsd.core system that the given signal happened. You can call this from inside a signal handler. 4944 +/ 4945 public static void markSignalOccurred(int sigNumber) nothrow { 4946 import core.sys.posix.unistd; 4947 4948 if(sigNumber == SIGCHLD) 4949 volatileStore(&sigChildHappened, 1); 4950 if(sigNumber == SIGINT) 4951 volatileStore(&sigIntrHappened, 1); 4952 4953 version(Arsd_core_epoll) { 4954 ulong writeValue = 1; 4955 write(signalPipeFd, &writeValue, writeValue.sizeof); 4956 } 4957 } 4958 } 4959 4960 version(Arsd_core_epoll) { 4961 4962 import core.sys.linux.epoll; 4963 import core.sys.linux.sys.eventfd; 4964 4965 private this() { 4966 4967 if(!globalsInitialized) { 4968 synchronized { 4969 if(!globalsInitialized) { 4970 // blocking signals is problematic because it is inherited by child processes 4971 // and that can be problematic for general purpose stuff so i use a self pipe 4972 // here. though since it is linux, im using an eventfd instead just to notify 4973 signalPipeFd = ErrnoEnforce!eventfd(0, EFD_CLOEXEC | EFD_NONBLOCK); 4974 signalReaderCallback = new CallbackHelper(&signalReader); 4975 4976 runInTaskRunnerQueue = new CallbackQueue("task runners", true); 4977 runInHelperThreadQueue = new CallbackQueue("helper threads", true); 4978 4979 setSignalHandlers(); 4980 4981 globalsInitialized = true; 4982 } 4983 } 4984 } 4985 4986 epollfd = epoll_create1(EPOLL_CLOEXEC); 4987 4988 // FIXME: ensure UI events get top priority 4989 4990 // global listeners 4991 4992 // FIXME: i should prolly keep the tokens and release them when tearing down. 4993 4994 cast(void) addCallbackOnFdReadable(signalPipeFd, signalReaderCallback); 4995 if(true) { // FIXME: if this is a task runner vs helper thread vs ui thread 4996 cast(void) addCallbackOnFdReadable(runInTaskRunnerQueue.fd, runInTaskRunnerQueue.callback); 4997 runInTaskRunnerQueue.callback.addref(); 4998 } else { 4999 cast(void) addCallbackOnFdReadable(runInHelperThreadQueue.fd, runInHelperThreadQueue.callback); 5000 runInHelperThreadQueue.callback.addref(); 5001 } 5002 5003 // local listener 5004 thisThreadQueue = new CallbackQueue("this thread", false); 5005 cast(void) addCallbackOnFdReadable(thisThreadQueue.fd, thisThreadQueue.callback); 5006 5007 // what are we going to do about timers? 5008 } 5009 5010 void teardown() { 5011 import core.sys.posix.fcntl; 5012 import core.sys.posix.unistd; 5013 5014 close(epollfd); 5015 epollfd = -1; 5016 5017 thisThreadQueue.teardown(); 5018 5019 // FIXME: should prolly free anything left in the callback queue, tho those could also be GC managed tbh. 5020 } 5021 5022 /+ // i call it explicitly at the thread exit instead, but worker threads aren't really supposed to exit generally speaking till process done anyway 5023 static ~this() { 5024 teardown(); 5025 } 5026 +/ 5027 5028 static void teardownGlobals() { 5029 import core.sys.posix.fcntl; 5030 import core.sys.posix.unistd; 5031 5032 synchronized { 5033 restoreSignalHandlers(); 5034 close(signalPipeFd); 5035 signalReaderCallback.release(); 5036 5037 runInTaskRunnerQueue.teardown(); 5038 runInHelperThreadQueue.teardown(); 5039 5040 globalsInitialized = false; 5041 } 5042 5043 } 5044 5045 5046 private static final class CallbackQueue { 5047 int fd = -1; 5048 string name; 5049 CallbackHelper callback; 5050 SynchronizedCircularBuffer!CallbackHelper queue; 5051 5052 this(string name, bool dequeueIsShared) { 5053 this.name = name; 5054 queue = typeof(queue)(this); 5055 5056 fd = ErrnoEnforce!eventfd(0, EFD_CLOEXEC | EFD_NONBLOCK | (dequeueIsShared ? EFD_SEMAPHORE : 0)); 5057 5058 callback = new CallbackHelper(dequeueIsShared ? &sharedDequeueCb : &threadLocalDequeueCb); 5059 } 5060 5061 bool resetEvent() { 5062 import core.sys.posix.unistd; 5063 ulong count; 5064 return read(fd, &count, count.sizeof) == count.sizeof; 5065 } 5066 5067 void sharedDequeueCb() { 5068 if(resetEvent()) { 5069 auto cb = queue.dequeue(); 5070 cb.call(); 5071 cb.release(); 5072 } 5073 } 5074 5075 void threadLocalDequeueCb() { 5076 CallbackHelper[16] buffer; 5077 foreach(cb; queue.dequeueSeveral(buffer[], () { resetEvent(); })) { 5078 cb.call(); 5079 cb.release(); 5080 } 5081 } 5082 5083 void enqueue(CallbackHelper cb) { 5084 if(queue.enqueue(cb)) { 5085 import core.sys.posix.unistd; 5086 ulong count = 1; 5087 ErrnoEnforce!write(fd, &count, count.sizeof); 5088 } else { 5089 throw new ArsdException!"queue is full"(name); 5090 } 5091 } 5092 5093 void teardown() { 5094 import core.sys.posix.fcntl; 5095 import core.sys.posix.unistd; 5096 5097 close(fd); 5098 fd = -1; 5099 5100 callback.release(); 5101 } 5102 } 5103 5104 // there's a global instance of this we refer back to 5105 private __gshared { 5106 bool globalsInitialized; 5107 5108 CallbackHelper signalReaderCallback; 5109 5110 CallbackQueue runInTaskRunnerQueue; 5111 CallbackQueue runInHelperThreadQueue; 5112 5113 int exitEventFd = -1; // FIXME: implement 5114 } 5115 5116 // and then the local loop 5117 private { 5118 int epollfd = -1; 5119 5120 CallbackQueue thisThreadQueue; 5121 } 5122 5123 // signal stuff { 5124 import core.sys.posix.signal; 5125 5126 private __gshared sigaction_t oldSigIntr; 5127 private __gshared sigaction_t oldSigChld; 5128 private __gshared sigaction_t oldSigPipe; 5129 5130 private __gshared int signalPipeFd = -1; 5131 // sigpipe not important, i handle errors on the writes 5132 5133 public static void setSignalHandlers() { 5134 static extern(C) void interruptHandler(int sigNumber) nothrow { 5135 markSignalOccurred(sigNumber); 5136 5137 /+ 5138 // calling the old handler is non-trivial since there can be ignore 5139 // or default or a plain handler or a sigaction 3 arg handler and i 5140 // i don't think it is worth teh complication 5141 sigaction_t* oldHandler; 5142 if(sigNumber == SIGCHLD) 5143 oldHandler = &oldSigChld; 5144 else if(sigNumber == SIGINT) 5145 oldHandler = &oldSigIntr; 5146 if(oldHandler && oldHandler.sa_handler) 5147 oldHandler 5148 +/ 5149 } 5150 5151 sigaction_t n; 5152 n.sa_handler = &interruptHandler; 5153 n.sa_mask = cast(sigset_t) 0; 5154 n.sa_flags = 0; 5155 sigaction(SIGINT, &n, &oldSigIntr); 5156 sigaction(SIGCHLD, &n, &oldSigChld); 5157 5158 n.sa_handler = SIG_IGN; 5159 sigaction(SIGPIPE, &n, &oldSigPipe); 5160 } 5161 5162 public static void restoreSignalHandlers() { 5163 sigaction(SIGINT, &oldSigIntr, null); 5164 sigaction(SIGCHLD, &oldSigChld, null); 5165 sigaction(SIGPIPE, &oldSigPipe, null); 5166 } 5167 5168 private static void signalReader() { 5169 import core.sys.posix.unistd; 5170 ulong number; 5171 read(signalPipeFd, &number, number.sizeof); 5172 5173 signalChecker(); 5174 } 5175 // signal stuff done } 5176 5177 // the any thread poll is just registered in the this thread poll w/ exclusive. nobody actaully epoll_waits 5178 // on the global one directly. 5179 5180 void runOnce() { 5181 epoll_event[16] events; 5182 auto ret = epoll_wait(epollfd, events.ptr, cast(int) events.length, -1); // FIXME: timeout 5183 if(ret == -1) { 5184 import core.stdc.errno; 5185 if(errno == EINTR) { 5186 return; 5187 } 5188 throw new ErrnoApiException("epoll_wait", errno); 5189 } else if(ret == 0) { 5190 // timeout 5191 } else { 5192 // loop events and call associated callbacks 5193 foreach(event; events[0 .. ret]) { 5194 auto flags = event.events; 5195 auto cbObject = cast(CallbackHelper) event.data.ptr; 5196 5197 // FIXME: or if it is an error... 5198 // EPOLLERR - write end of pipe when read end closed or other error. and EPOLLHUP - terminal hangup or read end when write end close (but it will give 0 reading after that soon anyway) 5199 5200 cbObject.call(); 5201 } 5202 } 5203 } 5204 5205 // building blocks for low-level integration with the loop 5206 5207 UnregisterToken addCallbackOnFdReadable(int fd, CallbackHelper cb) { 5208 epoll_event event; 5209 event.data.ptr = cast(void*) cb; 5210 event.events = EPOLLIN | EPOLLEXCLUSIVE; 5211 if(epoll_ctl(epollfd, EPOLL_CTL_ADD, fd, &event) == -1) 5212 throw new ErrnoApiException("epoll_ctl", errno); 5213 5214 return UnregisterToken(this, fd, cb); 5215 } 5216 5217 /++ 5218 Adds a one-off callback that you can optionally rearm when it happens. 5219 +/ 5220 RearmToken addCallbackOnFdReadableOneShot(int fd, CallbackHelper cb) { 5221 epoll_event event; 5222 event.data.ptr = cast(void*) cb; 5223 event.events = EPOLLIN | EPOLLONESHOT; 5224 if(epoll_ctl(epollfd, EPOLL_CTL_ADD, fd, &event) == -1) 5225 throw new ErrnoApiException("epoll_ctl", errno); 5226 5227 return RearmToken(true, this, fd, cb, EPOLLIN | EPOLLONESHOT); 5228 } 5229 5230 /++ 5231 Adds a one-off callback that you can optionally rearm when it happens. 5232 +/ 5233 RearmToken addCallbackOnFdWritableOneShot(int fd, CallbackHelper cb) { 5234 epoll_event event; 5235 event.data.ptr = cast(void*) cb; 5236 event.events = EPOLLOUT | EPOLLONESHOT; 5237 if(epoll_ctl(epollfd, EPOLL_CTL_ADD, fd, &event) == -1) 5238 throw new ErrnoApiException("epoll_ctl", errno); 5239 5240 return RearmToken(false, this, fd, cb, EPOLLOUT | EPOLLONESHOT); 5241 } 5242 5243 private void unregisterFd(int fd) { 5244 epoll_event event; 5245 if(epoll_ctl(epollfd, EPOLL_CTL_DEL, fd, &event) == -1) 5246 throw new ErrnoApiException("epoll_ctl", errno); 5247 } 5248 5249 private void rearmFd(RearmToken token) { 5250 epoll_event event; 5251 event.data.ptr = cast(void*) token.cb; 5252 event.events = token.flags; 5253 if(epoll_ctl(epollfd, EPOLL_CTL_MOD, token.fd, &event) == -1) 5254 throw new ErrnoApiException("epoll_ctl", errno); 5255 } 5256 5257 // Disk files will have to be sent as messages to a worker to do the read and report back a completion packet. 5258 } 5259 5260 version(Arsd_core_kqueue) { 5261 // FIXME 5262 } 5263 5264 // cross platform adapters 5265 void setTimeout() {} 5266 void addFileOrDirectoryChangeListener(FilePath name, uint flags, bool recursive = false) {} 5267 } 5268 5269 // deduplication???????// 5270 bool postMessage(ThreadToRunIn destination, void delegate() code) { 5271 return false; 5272 } 5273 bool postMessage(ThreadToRunIn destination, Object message) { 5274 return false; 5275 } 5276 5277 /+ 5278 void main() { 5279 // FIXME: the offset doesn't seem to be done right 5280 auto file = new AsyncFile(FilePath("test.txt"), AsyncFile.OpenMode.writeWithTruncation); 5281 file.write("hello", 10).waitForCompletion(); 5282 } 5283 +/ 5284 5285 // to test the mailboxes 5286 /+ 5287 void main() { 5288 /+ 5289 import std.stdio; 5290 Thread[4] pool; 5291 5292 bool shouldExit; 5293 5294 static int received; 5295 5296 static void tester() { 5297 received++; 5298 //writeln(cast(void*) Thread.getThis, " ", received); 5299 } 5300 5301 foreach(ref thread; pool) { 5302 thread = new Thread(() { 5303 getThisThreadEventLoop().run(() { 5304 return shouldExit; 5305 }); 5306 }); 5307 thread.start(); 5308 } 5309 5310 getThisThreadEventLoop(); // ensure it is all initialized before proceeding. FIXME: i should have an ensure initialized function i do on most the public apis. 5311 5312 int lol; 5313 5314 try 5315 foreach(i; 0 .. 6000) { 5316 CoreEventLoopImplementation.runInTaskRunnerQueue.enqueue(new CallbackHelper(&tester)); 5317 lol = cast(int) i; 5318 } 5319 catch(ArsdExceptionBase e) { 5320 Thread.sleep(50.msecs); 5321 writeln(e); 5322 writeln(lol); 5323 } 5324 5325 import core.stdc.stdlib; 5326 exit(0); 5327 5328 version(none) 5329 foreach(i; 0 .. 100) 5330 CoreEventLoopImplementation.runInTaskRunnerQueue.enqueue(new CallbackHelper(&tester)); 5331 5332 5333 foreach(ref thread; pool) { 5334 thread.join(); 5335 } 5336 +/ 5337 5338 5339 static int received; 5340 5341 static void tester() { 5342 received++; 5343 //writeln(cast(void*) Thread.getThis, " ", received); 5344 } 5345 5346 5347 5348 auto ev = cast(CoreEventLoopImplementation) getThisThreadEventLoop(); 5349 foreach(i; 0 .. 100) 5350 ev.thisThreadQueue.enqueue(new CallbackHelper(&tester)); 5351 foreach(i; 0 .. 100 / 16 + 1) 5352 ev.runOnce(); 5353 import std.conv; 5354 assert(received == 100, to!string(received)); 5355 5356 } 5357 +/ 5358 5359 /++ 5360 This is primarily a helper for the event queues. It is public in the hope it might be useful, 5361 but subject to change without notice; I will treat breaking it the same as if it is private. 5362 (That said, it is a simple little utility that does its job, so it is unlikely to change much. 5363 The biggest change would probably be letting it grow and changing from inline to dynamic array.) 5364 5365 It is a fixed-size ring buffer that synchronizes on a given object you give it in the constructor. 5366 5367 After enqueuing something, you should probably set an event to notify the other threads. This is left 5368 as an exercise to you (or another wrapper). 5369 +/ 5370 struct SynchronizedCircularBuffer(T, size_t maxSize = 128) { 5371 private T[maxSize] ring; 5372 private int front; 5373 private int back; 5374 5375 private Object synchronizedOn; 5376 5377 @disable this(); 5378 5379 /++ 5380 The Object's monitor is used to synchronize the methods in here. 5381 +/ 5382 this(Object synchronizedOn) { 5383 this.synchronizedOn = synchronizedOn; 5384 } 5385 5386 /++ 5387 Note the potential race condition between calling this and actually dequeuing something. You might 5388 want to acquire the lock on the object before calling this (nested synchronized things are allowed 5389 as long as the same thread is the one doing it). 5390 +/ 5391 bool isEmpty() { 5392 synchronized(this.synchronizedOn) { 5393 return front == back; 5394 } 5395 } 5396 5397 /++ 5398 Note the potential race condition between calling this and actually queuing something. 5399 +/ 5400 bool isFull() { 5401 synchronized(this.synchronizedOn) { 5402 return isFullUnsynchronized(); 5403 } 5404 } 5405 5406 private bool isFullUnsynchronized() nothrow const { 5407 return ((back + 1) % ring.length) == front; 5408 5409 } 5410 5411 /++ 5412 If this returns true, you should signal listening threads (with an event or a semaphore, 5413 depending on how you dequeue it). If it returns false, the queue was full and your thing 5414 was NOT added. You might wait and retry later (you could set up another event to signal it 5415 has been read and wait for that, or maybe try on a timer), or just fail and throw an exception 5416 or to abandon the message. 5417 +/ 5418 bool enqueue(T what) { 5419 synchronized(this.synchronizedOn) { 5420 if(isFullUnsynchronized()) 5421 return false; 5422 ring[(back++) % ring.length] = what; 5423 return true; 5424 } 5425 } 5426 5427 private T dequeueUnsynchronized() nothrow { 5428 assert(front != back); 5429 return ring[(front++) % ring.length]; 5430 } 5431 5432 /++ 5433 If you are using a semaphore to signal, you can call this once for each count of it 5434 and you can do that separately from this call (though they should be paired). 5435 5436 If you are using an event, you should use [dequeueSeveral] instead to drain it. 5437 +/ 5438 T dequeue() { 5439 synchronized(this.synchronizedOn) { 5440 return dequeueUnsynchronized(); 5441 } 5442 } 5443 5444 /++ 5445 Note that if you use a semaphore to signal waiting threads, you should probably not call this. 5446 5447 If you use a set/reset event, there's a potential race condition between the dequeue and event 5448 reset. This is why the `runInsideLockIfEmpty` delegate is there - when it is empty, before it 5449 unlocks, it will give you a chance to reset the event. Otherwise, it can remain set to indicate 5450 that there's still pending data in the queue. 5451 +/ 5452 T[] dequeueSeveral(return T[] buffer, scope void delegate() runInsideLockIfEmpty = null) { 5453 int pos; 5454 synchronized(this.synchronizedOn) { 5455 while(pos < buffer.length && front != back) { 5456 buffer[pos++] = dequeueUnsynchronized(); 5457 } 5458 if(front == back && runInsideLockIfEmpty !is null) 5459 runInsideLockIfEmpty(); 5460 } 5461 return buffer[0 .. pos]; 5462 } 5463 } 5464 5465 unittest { 5466 Object object = new Object(); 5467 auto queue = SynchronizedCircularBuffer!CallbackHelper(object); 5468 assert(queue.isEmpty); 5469 foreach(i; 0 .. queue.ring.length - 1) 5470 queue.enqueue(cast(CallbackHelper) cast(void*) i); 5471 assert(queue.isFull); 5472 5473 foreach(i; 0 .. queue.ring.length - 1) 5474 assert(queue.dequeue() is (cast(CallbackHelper) cast(void*) i)); 5475 assert(queue.isEmpty); 5476 5477 foreach(i; 0 .. queue.ring.length - 1) 5478 queue.enqueue(cast(CallbackHelper) cast(void*) i); 5479 assert(queue.isFull); 5480 5481 CallbackHelper[] buffer = new CallbackHelper[](300); 5482 auto got = queue.dequeueSeveral(buffer); 5483 assert(got.length == queue.ring.length - 1); 5484 assert(queue.isEmpty); 5485 foreach(i, item; got) 5486 assert(item is (cast(CallbackHelper) cast(void*) i)); 5487 5488 foreach(i; 0 .. 8) 5489 queue.enqueue(cast(CallbackHelper) cast(void*) i); 5490 buffer = new CallbackHelper[](4); 5491 got = queue.dequeueSeveral(buffer); 5492 assert(got.length == 4); 5493 foreach(i, item; got) 5494 assert(item is (cast(CallbackHelper) cast(void*) i)); 5495 got = queue.dequeueSeveral(buffer); 5496 assert(got.length == 4); 5497 foreach(i, item; got) 5498 assert(item is (cast(CallbackHelper) cast(void*) (i+4))); 5499 got = queue.dequeueSeveral(buffer); 5500 assert(got.length == 0); 5501 assert(queue.isEmpty); 5502 } 5503 5504 /++ 5505 5506 +/ 5507 enum ByteOrder { 5508 irrelevant, 5509 littleEndian, 5510 bigEndian, 5511 } 5512 5513 /++ 5514 A class to help write a stream of binary data to some target. 5515 5516 NOT YET FUNCTIONAL 5517 +/ 5518 class WritableStream { 5519 /++ 5520 5521 +/ 5522 this(size_t bufferSize) { 5523 this(new ubyte[](bufferSize)); 5524 } 5525 5526 /// ditto 5527 this(ubyte[] buffer) { 5528 this.buffer = buffer; 5529 } 5530 5531 /++ 5532 5533 +/ 5534 final void put(T)(T value, ByteOrder byteOrder = ByteOrder.irrelevant, string file = __FILE__, size_t line = __LINE__) { 5535 static if(T.sizeof == 8) 5536 ulong b; 5537 else static if(T.sizeof == 4) 5538 uint b; 5539 else static if(T.sizeof == 2) 5540 ushort b; 5541 else static if(T.sizeof == 1) 5542 ubyte b; 5543 else static assert(0, "unimplemented type, try using just the basic types"); 5544 5545 if(byteOrder == ByteOrder.irrelevant && T.sizeof > 1) 5546 throw new InvalidArgumentsException("byteOrder", "byte order must be specified for type " ~ T.stringof ~ " because it is bigger than one byte", "WritableStream.put", file, line); 5547 5548 final switch(byteOrder) { 5549 case ByteOrder.irrelevant: 5550 writeOneByte(b); 5551 break; 5552 case ByteOrder.littleEndian: 5553 foreach(i; 0 .. T.sizeof) { 5554 writeOneByte(b & 0xff); 5555 b >>= 8; 5556 } 5557 break; 5558 case ByteOrder.bigEndian: 5559 int amount = T.sizeof * 8 - 8; 5560 foreach(i; 0 .. T.sizeof) { 5561 writeOneByte((b >> amount) & 0xff); 5562 amount -= 8; 5563 } 5564 break; 5565 } 5566 } 5567 5568 /// ditto 5569 final void put(T : E[], E)(T value, ByteOrder elementByteOrder = ByteOrder.irrelevant, string file = __FILE__, size_t line = __LINE__) { 5570 foreach(item; value) 5571 put(item, elementByteOrder, file, line); 5572 } 5573 5574 /++ 5575 Performs a final flush() call, then marks the stream as closed, meaning no further data will be written to it. 5576 +/ 5577 void close() { 5578 isClosed_ = true; 5579 } 5580 5581 /++ 5582 Writes what is currently in the buffer to the target and waits for the target to accept it. 5583 Please note: if you are subclassing this to go to a different target 5584 +/ 5585 void flush() {} 5586 5587 /++ 5588 Returns true if either you closed it or if the receiving end closed their side, indicating they 5589 don't want any more data. 5590 +/ 5591 bool isClosed() { 5592 return isClosed_; 5593 } 5594 5595 // hasRoomInBuffer 5596 // canFlush 5597 // waitUntilCanFlush 5598 5599 // flushImpl 5600 // markFinished / close - tells the other end you're done 5601 5602 private final writeOneByte(ubyte value) { 5603 if(bufferPosition == buffer.length) 5604 flush(); 5605 5606 buffer[bufferPosition++] = value; 5607 } 5608 5609 5610 private { 5611 ubyte[] buffer; 5612 int bufferPosition; 5613 bool isClosed_; 5614 } 5615 } 5616 5617 /++ 5618 A stream can be used by just one task at a time, but one task can consume multiple streams. 5619 5620 Streams may be populated by async sources (in which case they must be called from a fiber task), 5621 from a function generating the data on demand (including an input range), from memory, or from a synchronous file. 5622 5623 A stream of heterogeneous types is compatible with input ranges. 5624 5625 It reads binary data. 5626 +/ 5627 version(HasThread) class ReadableStream { 5628 5629 this() { 5630 5631 } 5632 5633 /++ 5634 Gets data of the specified type `T` off the stream. The byte order of the T on the stream must be specified unless it is irrelevant (e.g. single byte entries). 5635 5636 --- 5637 // get an int out of a big endian stream 5638 int i = stream.get!int(ByteOrder.bigEndian); 5639 5640 // get i bytes off the stream 5641 ubyte[] data = stream.get!(ubyte[])(i); 5642 --- 5643 +/ 5644 final T get(T)(ByteOrder byteOrder = ByteOrder.irrelevant, string file = __FILE__, size_t line = __LINE__) { 5645 if(byteOrder == ByteOrder.irrelevant && T.sizeof > 1) 5646 throw new InvalidArgumentsException("byteOrder", "byte order must be specified for type " ~ T.stringof ~ " because it is bigger than one byte", "ReadableStream.get", file, line); 5647 5648 // FIXME: what if it is a struct? 5649 5650 while(bufferedLength() < T.sizeof) 5651 waitForAdditionalData(); 5652 5653 static if(T.sizeof == 1) { 5654 ubyte ret = consumeOneByte(); 5655 return *cast(T*) &ret; 5656 } else { 5657 static if(T.sizeof == 8) 5658 ulong ret; 5659 else static if(T.sizeof == 4) 5660 uint ret; 5661 else static if(T.sizeof == 2) 5662 ushort ret; 5663 else static assert(0, "unimplemented type, try using just the basic types"); 5664 5665 if(byteOrder == ByteOrder.littleEndian) { 5666 typeof(ret) buffer; 5667 foreach(b; 0 .. T.sizeof) { 5668 buffer = consumeOneByte(); 5669 buffer <<= T.sizeof * 8 - 8; 5670 5671 ret >>= 8; 5672 ret |= buffer; 5673 } 5674 } else { 5675 foreach(b; 0 .. T.sizeof) { 5676 ret <<= 8; 5677 ret |= consumeOneByte(); 5678 } 5679 } 5680 5681 return *cast(T*) &ret; 5682 } 5683 } 5684 5685 /// ditto 5686 final T get(T : E[], E)(size_t length, ByteOrder elementByteOrder = ByteOrder.irrelevant, string file = __FILE__, size_t line = __LINE__) { 5687 if(elementByteOrder == ByteOrder.irrelevant && E.sizeof > 1) 5688 throw new InvalidArgumentsException("elementByteOrder", "byte order must be specified for type " ~ E.stringof ~ " because it is bigger than one byte", "ReadableStream.get", file, line); 5689 5690 // if the stream is closed before getting the length or the terminator, should we send partial stuff 5691 // or just throw? 5692 5693 while(bufferedLength() < length * E.sizeof) 5694 waitForAdditionalData(); 5695 5696 T ret; 5697 5698 ret.length = length; 5699 5700 if(false && elementByteOrder == ByteOrder.irrelevant) { 5701 // ret[] = 5702 // FIXME: can prolly optimize 5703 } else { 5704 foreach(i; 0 .. length) 5705 ret[i] = get!E(elementByteOrder); 5706 } 5707 5708 return ret; 5709 5710 } 5711 5712 /// ditto 5713 final T get(T : E[], E)(scope bool delegate(E e) isTerminatingSentinel, ByteOrder elementByteOrder = ByteOrder.irrelevant, string file = __FILE__, size_t line = __LINE__) { 5714 if(byteOrder == ByteOrder.irrelevant && E.sizeof > 1) 5715 throw new InvalidArgumentsException("elementByteOrder", "byte order must be specified for type " ~ E.stringof ~ " because it is bigger than one byte", "ReadableStream.get", file, line); 5716 5717 assert(0, "Not implemented"); 5718 } 5719 5720 /++ 5721 5722 +/ 5723 bool isClosed() { 5724 return isClosed_; 5725 } 5726 5727 // Control side of things 5728 5729 private bool isClosed_; 5730 5731 /++ 5732 Feeds data into the stream, which can be consumed by `get`. If a task is waiting for more 5733 data to satisfy its get requests, this will trigger those tasks to resume. 5734 5735 If you feed it empty data, it will mark the stream as closed. 5736 +/ 5737 void feedData(ubyte[] data) { 5738 if(data.length == 0) 5739 isClosed_ = true; 5740 5741 currentBuffer = data; 5742 // this is a borrowed buffer, so we won't keep the reference long term 5743 scope(exit) 5744 currentBuffer = null; 5745 5746 if(waitingTask !is null) { 5747 waitingTask.call(); 5748 } 5749 } 5750 5751 /++ 5752 You basically have to use this thing from a task 5753 +/ 5754 protected void waitForAdditionalData() { 5755 Fiber task = Fiber.getThis; 5756 5757 assert(task !is null); 5758 5759 if(waitingTask !is null && waitingTask !is task) 5760 throw new ArsdException!"streams can only have one waiting task"; 5761 5762 // copy any pending data in our buffer to the longer-term buffer 5763 if(currentBuffer.length) 5764 leftoverBuffer ~= currentBuffer; 5765 5766 waitingTask = task; 5767 task.yield(); 5768 } 5769 5770 private Fiber waitingTask; 5771 private ubyte[] leftoverBuffer; 5772 private ubyte[] currentBuffer; 5773 5774 private size_t bufferedLength() { 5775 return leftoverBuffer.length + currentBuffer.length; 5776 } 5777 5778 private ubyte consumeOneByte() { 5779 ubyte b; 5780 if(leftoverBuffer.length) { 5781 b = leftoverBuffer[0]; 5782 leftoverBuffer = leftoverBuffer[1 .. $]; 5783 } else if(currentBuffer.length) { 5784 b = currentBuffer[0]; 5785 currentBuffer = currentBuffer[1 .. $]; 5786 } else { 5787 assert(0, "consuming off an empty buffer is impossible"); 5788 } 5789 5790 return b; 5791 } 5792 } 5793 5794 // FIXME: do a stringstream too 5795 5796 unittest { 5797 auto stream = new ReadableStream(); 5798 5799 int position; 5800 char[16] errorBuffer; 5801 5802 auto fiber = new Fiber(() { 5803 position = 1; 5804 int a = stream.get!int(ByteOrder.littleEndian); 5805 assert(a == 10, intToString(a, errorBuffer[])); 5806 position = 2; 5807 ubyte b = stream.get!ubyte; 5808 assert(b == 33); 5809 position = 3; 5810 5811 // ubyte[] c = stream.get!(ubyte[])(3); 5812 // int[] d = stream.get!(int[])(3); 5813 }); 5814 5815 fiber.call(); 5816 assert(position == 1); 5817 stream.feedData([10, 0, 0, 0]); 5818 assert(position == 2); 5819 stream.feedData([33]); 5820 assert(position == 3); 5821 5822 // stream.feedData([1,2,3]); 5823 // stream.feedData([1,2,3,4,1,2,3,4,1,2,3,4]); 5824 } 5825 5826 /++ 5827 UNSTABLE, NOT FULLY IMPLEMENTED. DO NOT USE YET. 5828 5829 You might use this like: 5830 5831 --- 5832 auto proc = new ExternalProcess(); 5833 auto stdoutStream = new ReadableStream(); 5834 5835 // to use a stream you can make one and have a task consume it 5836 runTask({ 5837 while(!stdoutStream.isClosed) { 5838 auto line = stdoutStream.get!string(e => e == '\n'); 5839 } 5840 }); 5841 5842 // then make the process feed into the stream 5843 proc.onStdoutAvailable = (got) { 5844 stdoutStream.feedData(got); // send it to the stream for processing 5845 stdout.rawWrite(got); // forward it through to our own thing 5846 // could also append it to a buffer to return it on complete 5847 }; 5848 proc.start(); 5849 --- 5850 5851 Please note that this does not currently and I have no plans as of this writing to add support for any kind of direct file descriptor passing. It always pipes them back to the parent for processing. If you don't want this, call the lower level functions yourself; the reason this class is here is to aid integration in the arsd.core event loop. Of course, I might change my mind on this. 5852 5853 Bugs: 5854 Not implemented at all on Windows yet. 5855 +/ 5856 class ExternalProcess /*: AsyncOperationRequest*/ { 5857 5858 private static version(Posix) { 5859 __gshared ExternalProcess[pid_t] activeChildren; 5860 5861 void recordChildCreated(pid_t pid, ExternalProcess proc) { 5862 synchronized(typeid(ExternalProcess)) { 5863 activeChildren[pid] = proc; 5864 } 5865 } 5866 5867 void recordChildTerminated(pid_t pid, int status) { 5868 synchronized(typeid(ExternalProcess)) { 5869 if(pid in activeChildren) { 5870 auto ac = activeChildren[pid]; 5871 ac.completed = true; 5872 ac.status = status; 5873 activeChildren.remove(pid); 5874 } 5875 } 5876 } 5877 } 5878 5879 // FIXME: config to pass through a shell or not 5880 5881 /++ 5882 This is the native version for Windows. 5883 +/ 5884 this(string program, string commandLine) { 5885 version(Posix) { 5886 assert(0, "not implemented command line to posix args yet"); 5887 } 5888 else throw new NotYetImplementedException(); 5889 } 5890 5891 this(string commandLine) { 5892 version(Posix) { 5893 assert(0, "not implemented command line to posix args yet"); 5894 } 5895 else throw new NotYetImplementedException(); 5896 } 5897 5898 this(string[] args) { 5899 version(Posix) { 5900 this.program = FilePath(args[0]); 5901 this.args = args; 5902 } 5903 else throw new NotYetImplementedException(); 5904 } 5905 5906 /++ 5907 This is the native version for Posix. 5908 +/ 5909 this(FilePath program, string[] args) { 5910 version(Posix) { 5911 this.program = program; 5912 this.args = args; 5913 } 5914 else throw new NotYetImplementedException(); 5915 } 5916 5917 // you can modify these before calling start 5918 int stdoutBufferSize = 32 * 1024; 5919 int stderrBufferSize = 8 * 1024; 5920 5921 void start() { 5922 version(Posix) { 5923 int ret; 5924 5925 int[2] stdinPipes; 5926 ret = pipe(stdinPipes); 5927 if(ret == -1) 5928 throw new ErrnoApiException("stdin pipe", errno); 5929 5930 scope(failure) { 5931 close(stdinPipes[0]); 5932 close(stdinPipes[1]); 5933 } 5934 5935 stdinFd = stdinPipes[1]; 5936 5937 int[2] stdoutPipes; 5938 ret = pipe(stdoutPipes); 5939 if(ret == -1) 5940 throw new ErrnoApiException("stdout pipe", errno); 5941 5942 scope(failure) { 5943 close(stdoutPipes[0]); 5944 close(stdoutPipes[1]); 5945 } 5946 5947 stdoutFd = stdoutPipes[0]; 5948 5949 int[2] stderrPipes; 5950 ret = pipe(stderrPipes); 5951 if(ret == -1) 5952 throw new ErrnoApiException("stderr pipe", errno); 5953 5954 scope(failure) { 5955 close(stderrPipes[0]); 5956 close(stderrPipes[1]); 5957 } 5958 5959 stderrFd = stderrPipes[0]; 5960 5961 5962 int[2] errorReportPipes; 5963 ret = pipe(errorReportPipes); 5964 if(ret == -1) 5965 throw new ErrnoApiException("error reporting pipe", errno); 5966 5967 scope(failure) { 5968 close(errorReportPipes[0]); 5969 close(errorReportPipes[1]); 5970 } 5971 5972 setCloExec(errorReportPipes[0]); 5973 setCloExec(errorReportPipes[1]); 5974 5975 auto forkRet = fork(); 5976 if(forkRet == -1) 5977 throw new ErrnoApiException("fork", errno); 5978 5979 if(forkRet == 0) { 5980 // child side 5981 5982 // FIXME can we do more error checking that is actually useful here? 5983 // these operations are virtually guaranteed to succeed given the setup anyway. 5984 5985 // FIXME pty too 5986 5987 void fail(int step) { 5988 import core.stdc.errno; 5989 auto code = errno; 5990 5991 // report the info back to the parent then exit 5992 5993 int[2] msg = [step, code]; 5994 auto ret = write(errorReportPipes[1], msg.ptr, msg.sizeof); 5995 5996 // but if this fails there's not much we can do... 5997 5998 import core.stdc.stdlib; 5999 exit(1); 6000 } 6001 6002 // dup2 closes the fd it is replacing automatically 6003 dup2(stdinPipes[0], 0); 6004 dup2(stdoutPipes[1], 1); 6005 dup2(stderrPipes[1], 2); 6006 6007 // don't need either of the original pipe fds anymore 6008 close(stdinPipes[0]); 6009 close(stdinPipes[1]); 6010 close(stdoutPipes[0]); 6011 close(stdoutPipes[1]); 6012 close(stderrPipes[0]); 6013 close(stderrPipes[1]); 6014 6015 // the error reporting pipe will be closed upon exec since we set cloexec before fork 6016 // and everything else should have cloexec set too hopefully. 6017 6018 if(beforeExec) 6019 beforeExec(); 6020 6021 // i'm not sure that a fully-initialized druntime is still usable 6022 // after a fork(), so i'm gonna stick to the C lib in here. 6023 6024 const(char)* file = mallocedStringz(program.path).ptr; 6025 if(file is null) 6026 fail(1); 6027 const(char)*[] argv = mallocSlice!(const(char)*)(args.length + 1); 6028 if(argv is null) 6029 fail(2); 6030 foreach(idx, arg; args) { 6031 argv[idx] = mallocedStringz(args[idx]).ptr; 6032 if(argv[idx] is null) 6033 fail(3); 6034 } 6035 argv[args.length] = null; 6036 6037 auto rete = execvp/*e*/(file, argv.ptr/*, envp*/); 6038 if(rete == -1) { 6039 fail(4); 6040 } else { 6041 // unreachable code, exec never returns if it succeeds 6042 assert(0); 6043 } 6044 } else { 6045 pid = forkRet; 6046 6047 recordChildCreated(pid, this); 6048 6049 // close our copy of the write side of the error reporting pipe 6050 // so the read will immediately give eof when the fork closes it too 6051 ErrnoEnforce!close(errorReportPipes[1]); 6052 6053 int[2] msg; 6054 // this will block to wait for it to actually either start up or fail to exec (which should be near instant) 6055 auto val = read(errorReportPipes[0], msg.ptr, msg.sizeof); 6056 6057 if(val == -1) 6058 throw new ErrnoApiException("read error report", errno); 6059 6060 if(val == msg.sizeof) { 6061 // error happened 6062 // FIXME: keep the step part of the error report too 6063 throw new ErrnoApiException("exec", msg[1]); 6064 } else if(val == 0) { 6065 // pipe closed, meaning exec succeeded 6066 } else { 6067 assert(0); // never supposed to happen 6068 } 6069 6070 // set the ones we keep to close upon future execs 6071 // FIXME should i set NOBLOCK at this time too? prolly should 6072 setCloExec(stdinPipes[1]); 6073 setCloExec(stdoutPipes[0]); 6074 setCloExec(stderrPipes[0]); 6075 6076 // and close the others 6077 ErrnoEnforce!close(stdinPipes[0]); 6078 ErrnoEnforce!close(stdoutPipes[1]); 6079 ErrnoEnforce!close(stderrPipes[1]); 6080 6081 ErrnoEnforce!close(errorReportPipes[0]); 6082 6083 // and now register the ones we need to read with the event loop so it can call the callbacks 6084 // also need to listen to SIGCHLD to queue up the terminated callback. FIXME 6085 6086 stdoutUnregisterToken = getThisThreadEventLoop().addCallbackOnFdReadable(stdoutFd, new CallbackHelper(&stdoutReadable)); 6087 stderrUnregisterToken = getThisThreadEventLoop().addCallbackOnFdReadable(stderrFd, new CallbackHelper(&stderrReadable)); 6088 } 6089 } 6090 } 6091 6092 private version(Posix) { 6093 import core.sys.posix.unistd; 6094 import core.sys.posix.fcntl; 6095 6096 int stdinFd = -1; 6097 int stdoutFd = -1; 6098 int stderrFd = -1; 6099 6100 ICoreEventLoop.UnregisterToken stdoutUnregisterToken; 6101 ICoreEventLoop.UnregisterToken stderrUnregisterToken; 6102 6103 pid_t pid = -1; 6104 6105 public void delegate() beforeExec; 6106 6107 FilePath program; 6108 string[] args; 6109 6110 void stdoutReadable() { 6111 if(stdoutReadBuffer is null) 6112 stdoutReadBuffer = new ubyte[](stdoutBufferSize); 6113 auto ret = read(stdoutFd, stdoutReadBuffer.ptr, stdoutReadBuffer.length); 6114 if(ret == -1) 6115 throw new ErrnoApiException("read", errno); 6116 if(onStdoutAvailable) { 6117 onStdoutAvailable(stdoutReadBuffer[0 .. ret]); 6118 } 6119 6120 if(ret == 0) { 6121 stdoutUnregisterToken.unregister(); 6122 6123 close(stdoutFd); 6124 stdoutFd = -1; 6125 } 6126 } 6127 6128 void stderrReadable() { 6129 if(stderrReadBuffer is null) 6130 stderrReadBuffer = new ubyte[](stderrBufferSize); 6131 auto ret = read(stderrFd, stderrReadBuffer.ptr, stderrReadBuffer.length); 6132 if(ret == -1) 6133 throw new ErrnoApiException("read", errno); 6134 if(onStderrAvailable) { 6135 onStderrAvailable(stderrReadBuffer[0 .. ret]); 6136 } 6137 6138 if(ret == 0) { 6139 stderrUnregisterToken.unregister(); 6140 6141 close(stderrFd); 6142 stderrFd = -1; 6143 } 6144 } 6145 } 6146 6147 private ubyte[] stdoutReadBuffer; 6148 private ubyte[] stderrReadBuffer; 6149 6150 void waitForCompletion() { 6151 getThisThreadEventLoop().run(&this.isComplete); 6152 } 6153 6154 bool isComplete() { 6155 return completed; 6156 } 6157 6158 bool completed; 6159 int status = int.min; 6160 6161 /++ 6162 If blocking, it will block the current task until the write succeeds. 6163 6164 Write `null` as data to close the pipe. Once the pipe is closed, you must not try to write to it again. 6165 +/ 6166 void writeToStdin(in void[] data) { 6167 version(Posix) { 6168 if(data is null) { 6169 close(stdinFd); 6170 stdinFd = -1; 6171 } else { 6172 // FIXME: check the return value again and queue async writes 6173 auto ret = write(stdinFd, data.ptr, data.length); 6174 if(ret == -1) 6175 throw new ErrnoApiException("write", errno); 6176 } 6177 } 6178 6179 } 6180 6181 void delegate(ubyte[] got) onStdoutAvailable; 6182 void delegate(ubyte[] got) onStderrAvailable; 6183 void delegate(int code) onTermination; 6184 6185 // pty? 6186 } 6187 6188 // FIXME: comment this out 6189 /+ 6190 unittest { 6191 auto proc = new ExternalProcess(FilePath("/bin/cat"), ["/bin/cat"]); 6192 6193 getThisThreadEventLoop(); // initialize it 6194 6195 int c = 0; 6196 proc.onStdoutAvailable = delegate(ubyte[] got) { 6197 if(c == 0) 6198 assert(cast(string) got == "hello!"); 6199 else 6200 assert(got.length == 0); 6201 // import std.stdio; writeln(got); 6202 c++; 6203 }; 6204 6205 proc.start(); 6206 6207 assert(proc.pid != -1); 6208 6209 6210 import std.stdio; 6211 Thread[4] pool; 6212 6213 bool shouldExit; 6214 6215 static int received; 6216 6217 proc.writeToStdin("hello!"); 6218 proc.writeToStdin(null); // closes the pipe 6219 6220 proc.waitForCompletion(); 6221 6222 assert(proc.status == 0); 6223 6224 assert(c == 2); 6225 6226 // writeln("here"); 6227 } 6228 +/ 6229 6230 // to test the thundering herd on signal handling 6231 version(none) 6232 unittest { 6233 Thread[4] pool; 6234 foreach(ref thread; pool) { 6235 thread = new class Thread { 6236 this() { 6237 super({ 6238 int count; 6239 getThisThreadEventLoop().run(() { 6240 if(count > 4) return true; 6241 count++; 6242 return false; 6243 }); 6244 }); 6245 } 6246 }; 6247 thread.start(); 6248 } 6249 foreach(ref thread; pool) { 6250 thread.join(); 6251 } 6252 } 6253 6254 /+ 6255 ================= 6256 STDIO REPLACEMENT 6257 ================= 6258 +/ 6259 6260 private void appendToBuffer(ref char[] buffer, ref int pos, scope const(char)[] what) { 6261 auto required = pos + what.length; 6262 if(buffer.length < required) 6263 buffer.length = required; 6264 buffer[pos .. pos + what.length] = what[]; 6265 pos += what.length; 6266 } 6267 6268 private void appendToBuffer(ref char[] buffer, ref int pos, long what) { 6269 if(buffer.length < pos + 16) 6270 buffer.length = pos + 16; 6271 auto sliced = intToString(what, buffer[pos .. $]); 6272 pos += sliced.length; 6273 } 6274 6275 /++ 6276 A `writeln` that actually works, at least for some basic types. 6277 6278 It works correctly on Windows, using the correct functions to write unicode to the console. even allocating a console if needed. If the output has been redirected to a file or pipe, it writes UTF-8. 6279 6280 This always does text. See also WritableStream and WritableTextStream when they are implemented. 6281 +/ 6282 void writeln(T...)(T t) { 6283 char[256] bufferBacking; 6284 char[] buffer = bufferBacking[]; 6285 int pos; 6286 6287 foreach(arg; t) { 6288 static if(is(typeof(arg) : const char[])) { 6289 appendToBuffer(buffer, pos, arg); 6290 } else static if(is(typeof(arg) : stringz)) { 6291 appendToBuffer(buffer, pos, arg.borrow); 6292 } else static if(is(typeof(arg) : long)) { 6293 appendToBuffer(buffer, pos, arg); 6294 } else static if(is(typeof(arg.toString()) : const char[])) { 6295 appendToBuffer(buffer, pos, arg.toString()); 6296 } else { 6297 appendToBuffer(buffer, pos, "<" ~ typeof(arg).stringof ~ ">"); 6298 } 6299 } 6300 6301 appendToBuffer(buffer, pos, "\n"); 6302 6303 actuallyWriteToStdout(buffer[0 .. pos]); 6304 } 6305 6306 private void actuallyWriteToStdout(scope char[] buffer) @trusted { 6307 6308 version(UseStdioWriteln) 6309 { 6310 import std.stdio; 6311 writeln(buffer); 6312 } 6313 else version(Windows) { 6314 import core.sys.windows.wincon; 6315 6316 auto hStdOut = GetStdHandle(STD_OUTPUT_HANDLE); 6317 if(hStdOut == null || hStdOut == INVALID_HANDLE_VALUE) { 6318 AllocConsole(); 6319 hStdOut = GetStdHandle(STD_OUTPUT_HANDLE); 6320 } 6321 6322 if(GetFileType(hStdOut) == FILE_TYPE_CHAR) { 6323 wchar[256] wbuffer; 6324 auto toWrite = makeWindowsString(buffer, wbuffer, WindowsStringConversionFlags.convertNewLines); 6325 6326 DWORD written; 6327 WriteConsoleW(hStdOut, toWrite.ptr, cast(DWORD) toWrite.length, &written, null); 6328 } else { 6329 DWORD written; 6330 WriteFile(hStdOut, buffer.ptr, cast(DWORD) buffer.length, &written, null); 6331 } 6332 } else { 6333 import unix = core.sys.posix.unistd; 6334 unix.write(1, buffer.ptr, buffer.length); 6335 } 6336 } 6337 6338 /+ 6339 6340 STDIO 6341 6342 /++ 6343 Please note using this will create a compile-time dependency on [arsd.terminal] 6344 6345 6346 6347 so my writeln replacement: 6348 6349 1) if the std output handle is null, alloc one 6350 2) if it is a character device, write out the proper Unicode text. 6351 3) otherwise write out UTF-8.... maybe with a BOM but maybe not. it is tricky to know what the other end of a pipe expects... 6352 [8:15 AM] 6353 im actually tempted to make the write binary to stdout functions throw an exception if it is a character console / interactive terminal instead of letting you spam it right out 6354 [8:16 AM] 6355 of course you can still cheat by casting binary data to string and using the write string function (and this might be appropriate sometimes) but there kinda is a legit difference between a text output and a binary output device 6356 6357 Stdout can represent either 6358 6359 +/ 6360 void writeln(){} { 6361 6362 } 6363 6364 stderr? 6365 6366 /++ 6367 Please note using this will create a compile-time dependency on [arsd.terminal] 6368 6369 It can be called from a task. 6370 6371 It works correctly on Windows and is user friendly on Linux (using arsd.terminal.getline) 6372 while also working if stdin has been redirected (where arsd.terminal itself would throw) 6373 6374 6375 so say you run a program on an interactive terminal. the program tries to open the stdin binary stream 6376 6377 instead of throwing, the prompt could change to indicate the binary data is expected and you can feed it in either by typing it up,,,, or running some command like maybe <file.bin to have the library do what the shell would have done and feed that to the rest of the program 6378 6379 +/ 6380 string readln()() { 6381 6382 } 6383 6384 6385 // if using stdio as a binary output thing you can pretend it is a file w/ stream capability 6386 struct File { 6387 WritableStream ostream; 6388 ReadableStream istream; 6389 6390 ulong tell; 6391 void seek(ulong to) {} 6392 6393 void sync(); 6394 void close(); 6395 } 6396 6397 // these are a bit special because if it actually is an interactive character device, it might be different than other files and even different than other pipes. 6398 WritableStream stdoutStream() { return null; } 6399 WritableStream stderrStream() { return null; } 6400 ReadableStream stdinStream() { return null; } 6401 6402 +/ 6403 6404 6405 /+ 6406 6407 6408 /+ 6409 Druntime appears to have stuff for darwin, freebsd. I might have to add some for openbsd here and maybe netbsd if i care to test it. 6410 +/ 6411 6412 /+ 6413 6414 arsd_core_init(number_of_worker_threads) 6415 6416 Building-block things wanted for the event loop integration: 6417 * ui 6418 * windows 6419 * terminal / console 6420 * generic 6421 * adopt fd 6422 * adopt windows handle 6423 * shared lib 6424 * load 6425 * timers (relative and real time) 6426 * create 6427 * update 6428 * cancel 6429 * file/directory watches 6430 * file created 6431 * file deleted 6432 * file modified 6433 * file ops 6434 * open 6435 * close 6436 * read 6437 * write 6438 * seek 6439 * sendfile on linux, TransmitFile on Windows 6440 * let completion handlers run in the io worker thread instead of signaling back 6441 * pipe ops (anonymous or named) 6442 * create 6443 * read 6444 * write 6445 * get info about other side of the pipe 6446 * network ops (stream + datagram, ip, ipv6, unix) 6447 * address look up 6448 * connect 6449 * start tls 6450 * listen 6451 * send 6452 * receive 6453 * get peer info 6454 * process ops 6455 * spawn 6456 * notifications when it is terminated or fork or execs 6457 * send signal 6458 * i/o pipes 6459 * thread ops (isDaemon?) 6460 * spawn 6461 * talk to its event loop 6462 * termination notification 6463 * signals 6464 * ctrl+c is the only one i really care about but the others might be made available too. sigchld needs to be done as an impl detail of process ops. 6465 * custom messages 6466 * should be able to send messages from finalizers... 6467 6468 * want to make sure i can stream stuff on top of it all too. 6469 6470 ======== 6471 6472 These things all refer back to a task-local thing that queues the tasks. If it is a fiber, it uses that 6473 and if it is a thread it uses that... 6474 6475 tls IArsdCoreEventLoop curentTaskInterface; // this yields on the wait for calls. the fiber swapper will swap this too. 6476 tls IArsdCoreEventLoop currentThreadInterface; // this blocks on the event loop 6477 6478 shared IArsdCoreEventLoop currentProcessInterface; // this dispatches to any available thread 6479 +/ 6480 6481 6482 /+ 6483 You might have configurable tasks that do not auto-start, e.g. httprequest. maybe @mustUse on those 6484 6485 then some that do auto-start, e.g. setTimeout 6486 6487 6488 timeouts: duration, MonoTime, or SysTime? duration is just a timer monotime auto-adjusts the when, systime sets a real time timerfd 6489 6490 tasks can be set to: 6491 thread affinity - this, any, specific reference 6492 reports to - defaults to this, can also pass down a parent reference. if reports to dies, all its subordinates are cancelled. 6493 6494 6495 you can send a message to a task... maybe maybe just to a task runner (which is itself a task?) 6496 6497 auto file = readFile(x); 6498 auto timeout = setTimeout(y); 6499 auto completed = waitForFirstToCompleteThenCancelOthers(file, timeout); 6500 if(completed == 0) { 6501 file.... 6502 } else { 6503 timeout.... 6504 } 6505 6506 /+ 6507 A task will run on a thread (with possible migration), and report to a task. 6508 +/ 6509 6510 // a compute task is run on a helper thread 6511 auto task = computeTask((shared(bool)* cancellationRequested) { 6512 // or pass in a yield thing... prolly a TaskController which has cancellationRequested and yield controls as well as send message to parent (sync or async) 6513 6514 // you'd periodically send messages back to the parent 6515 }, RunOn.AnyAvailable, Affinity.CanMigrate); 6516 6517 auto task = Task((TaskController controller) { 6518 foreach(x, 0 .. 1000) { 6519 if(x % 10 == 0) 6520 controller.yield(); // periodically yield control, which also checks for cancellation for us 6521 // do some work 6522 6523 controller.sendMessage(...); 6524 controller.sendProgress(x); // yields it for a foreach stream kind of thing 6525 } 6526 6527 return something; // automatically sends the something as the result in a TaskFinished message 6528 }); 6529 6530 foreach(item; task) // waitsForProgress, sendProgress sends an item and the final return sends an item 6531 {} 6532 6533 6534 see ~/test/task.d 6535 6536 // an io task is run locally via the event loops 6537 auto task2 = ioTask(() { 6538 6539 }); 6540 6541 6542 6543 waitForEvent 6544 +/ 6545 6546 /+ 6547 Most functions should prolly take a thread arg too, which defaults 6548 to this thread, but you can also pass it a reference, or a "any available" thing. 6549 6550 This can be a ufcs overload 6551 +/ 6552 6553 interface SemiSynchronousTask { 6554 6555 } 6556 6557 struct TimeoutCompletionResult { 6558 bool completed; 6559 6560 bool opCast(T : bool)() { 6561 return completed; 6562 } 6563 } 6564 6565 struct Timeout { 6566 void reschedule(Duration when) { 6567 6568 } 6569 6570 void cancel() { 6571 6572 } 6573 6574 TimeoutCompletionResult waitForCompletion() { 6575 return TimeoutCompletionResult(false); 6576 } 6577 } 6578 6579 Timeout setTimeout(void delegate() dg, int msecs, int permittedJitter = 20) { 6580 return Timeout.init; 6581 } 6582 6583 void clearTimeout(Timeout timeout) { 6584 timeout.cancel(); 6585 } 6586 6587 void createInterval() {} 6588 void clearInterval() {} 6589 6590 /++ 6591 Schedules a task at the given wall clock time. 6592 +/ 6593 void scheduleTask() {} 6594 6595 struct IoOperationCompletionResult { 6596 enum Status { 6597 cancelled, 6598 completed 6599 } 6600 6601 Status status; 6602 6603 int error; 6604 int bytesWritten; 6605 6606 bool opCast(T : bool)() { 6607 return status == Status.completed; 6608 } 6609 } 6610 6611 struct IoOperation { 6612 void cancel() {} 6613 6614 IoOperationCompletionResult waitForCompletion() { 6615 return IoOperationCompletionResult.init; 6616 } 6617 6618 // could contain a scoped class in here too so it stack allocated 6619 } 6620 6621 // Should return both the object and the index in the array! 6622 Result waitForFirstToComplete(Operation[]...) {} 6623 6624 IoOperation read(IoHandle handle, ubyte[] buffer 6625 6626 /+ 6627 class IoOperation {} 6628 6629 // an io operation and its buffer must not be modified or freed 6630 // in between a call to enqueue and a call to waitForCompletion 6631 // if you used the whenComplete callback, make sure it is NOT gc'd or scope thing goes out of scope in the mean time 6632 // if its dtor runs, it'd be forced to be cancelled... 6633 6634 scope IoOperation op = new IoOperation(buffer_size); 6635 op.start(); 6636 op.waitForCompletion(); 6637 +/ 6638 6639 /+ 6640 will want: 6641 read, write 6642 send, recv 6643 6644 cancel 6645 6646 open file, open (named or anonymous) pipe, open process 6647 connect, accept 6648 SSL 6649 close 6650 6651 postEvent 6652 postAPC? like run in gui thread / async 6653 waitForEvent ? needs to handle a timeout and a cancellation. would only work in the fiber task api. 6654 6655 waitForSuccess 6656 6657 interrupt handler 6658 6659 onPosixReadReadiness 6660 onPosixWriteReadiness 6661 6662 onWindowsHandleReadiness 6663 - but they're one-offs so you gotta reregister for each event 6664 +/ 6665 6666 6667 6668 /+ 6669 arsd.core.uda 6670 6671 you define a model struct with the types you want to extract 6672 6673 you get it with like Model extract(Model, UDAs...)(Model default) 6674 6675 defaultModel!alias > defaultModel!Type(defaultModel("identifier")) 6676 6677 6678 6679 6680 6681 6682 6683 6684 6685 6686 so while i laid there sleep deprived i did think a lil more on some uda stuff. it isn't especially novel but a combination of a few other techniques 6687 6688 you might be like 6689 6690 struct MyUdas { 6691 DbName name; 6692 DbIgnore ignore; 6693 } 6694 6695 elsewhere 6696 6697 foreach(alias; allMembers) { 6698 auto udas = getUdas!(MyUdas, __traits(getAttributes, alias))(MyUdas(DbName(__traits(identifier, alias)))); 6699 } 6700 6701 6702 so you pass the expected type and the attributes as the template params, then the runtime params are the default values for the given types 6703 6704 so what the thing does essentially is just sets the values of the given thing to the udas based on type then returns the modified instance 6705 6706 so the end result is you keep the last ones. it wouldn't report errors if multiple things added but it p simple to understand, simple to document (even though the default values are not in the struct itself, you can put ddocs in them), and uses the tricks to minimize generated code size 6707 +/ 6708 6709 +/ 6710 6711 package(arsd) version(Windows) extern(Windows) { 6712 BOOL CancelIoEx(HANDLE, LPOVERLAPPED); 6713 6714 struct WSABUF { 6715 ULONG len; 6716 ubyte* buf; 6717 } 6718 alias LPWSABUF = WSABUF*; 6719 6720 // https://learn.microsoft.com/en-us/windows/win32/api/winsock2/ns-winsock2-wsaoverlapped 6721 // "The WSAOVERLAPPED structure is compatible with the Windows OVERLAPPED structure." 6722 // so ima lie here in the bindings. 6723 6724 int WSASend(SOCKET, LPWSABUF, DWORD, LPDWORD, DWORD, LPOVERLAPPED, LPOVERLAPPED_COMPLETION_ROUTINE); 6725 int WSASendTo(SOCKET, LPWSABUF, DWORD, LPDWORD, DWORD, const sockaddr*, int, LPOVERLAPPED, LPOVERLAPPED_COMPLETION_ROUTINE); 6726 6727 int WSARecv(SOCKET, LPWSABUF, DWORD, LPDWORD, LPDWORD, LPOVERLAPPED, LPOVERLAPPED_COMPLETION_ROUTINE); 6728 int WSARecvFrom(SOCKET, LPWSABUF, DWORD, LPDWORD, LPDWORD, sockaddr*, LPINT, LPOVERLAPPED, LPOVERLAPPED_COMPLETION_ROUTINE); 6729 } 6730 6731 package(arsd) version(OSXCocoa) { 6732 6733 /+ 6734 To let Cocoa know that you intend to use multiple threads, all you have to do is spawn a single thread using the NSThread class and let that thread immediately exit. Your thread entry point need not do anything. Just the act of spawning a thread using NSThread is enough to ensure that the locks needed by the Cocoa frameworks are put in place. 6735 6736 If you are not sure if Cocoa thinks your application is multithreaded or not, you can use the isMultiThreaded method of NSThread to check. 6737 +/ 6738 6739 6740 struct DeifiedNSString { 6741 char[16] sso; 6742 const(char)[] str; 6743 6744 this(NSString s) { 6745 auto len = s.length; 6746 if(len <= sso.length / 4) 6747 str = sso[]; 6748 else 6749 str = new char[](len * 4); 6750 6751 NSUInteger count; 6752 NSRange leftover; 6753 auto ret = s.getBytes(cast(char*) str.ptr, str.length, &count, NSStringEncoding.NSUTF8StringEncoding, NSStringEncodingConversionOptions.none, NSRange(0, len), &leftover); 6754 if(ret) 6755 str = str[0 .. count]; 6756 else 6757 throw new Exception("uh oh"); 6758 } 6759 } 6760 6761 extern (Objective-C) { 6762 import core.attribute; // : selector, optional; 6763 6764 alias NSUInteger = size_t; 6765 alias NSInteger = ptrdiff_t; 6766 alias unichar = wchar; 6767 struct SEL_; 6768 alias SEL_* SEL; 6769 // this is called plain `id` in objective C but i fear mistakes with that in D. like sure it is a type instead of a variable like most things called id but i still think it is weird. i might change my mind later. 6770 alias void* NSid; // FIXME? the docs say this is a pointer to an instance of a class, but that is not necessary a child of NSObject 6771 6772 extern class NSObject { 6773 static NSObject alloc() @selector("alloc"); 6774 NSObject init() @selector("init"); 6775 6776 void retain() @selector("retain"); 6777 void release() @selector("release"); 6778 void autorelease() @selector("autorelease"); 6779 6780 void performSelectorOnMainThread(SEL aSelector, NSid arg, bool waitUntilDone) @selector("performSelectorOnMainThread:withObject:waitUntilDone:"); 6781 } 6782 6783 // this is some kind of generic in objc... 6784 extern class NSArray : NSObject { 6785 static NSArray arrayWithObjects(NSid* objects, NSUInteger count) @selector("arrayWithObjects:count:"); 6786 } 6787 6788 extern class NSString : NSObject { 6789 override static NSString alloc() @selector("alloc"); 6790 override NSString init() @selector("init"); 6791 6792 NSString initWithUTF8String(const scope char* str) @selector("initWithUTF8String:"); 6793 6794 NSString initWithBytes( 6795 const(ubyte)* bytes, 6796 NSUInteger length, 6797 NSStringEncoding encoding 6798 ) @selector("initWithBytes:length:encoding:"); 6799 6800 unichar characterAtIndex(NSUInteger index) @selector("characterAtIndex:"); 6801 NSUInteger length() @selector("length"); 6802 const char* UTF8String() @selector("UTF8String"); 6803 6804 void getCharacters(wchar* buffer, NSRange range) @selector("getCharacters:range:"); 6805 6806 bool getBytes(void* buffer, NSUInteger maxBufferCount, NSUInteger* usedBufferCount, NSStringEncoding encoding, NSStringEncodingConversionOptions options, NSRange range, NSRange* leftover) @selector("getBytes:maxLength:usedLength:encoding:options:range:remainingRange:"); 6807 } 6808 6809 struct NSRange { 6810 NSUInteger loc; 6811 NSUInteger len; 6812 } 6813 6814 enum NSStringEncodingConversionOptions : NSInteger { 6815 none = 0, 6816 NSAllowLossyEncodingConversion = 1, 6817 NSExternalRepresentationEncodingConversion = 2 6818 } 6819 6820 enum NSEventType { 6821 idk 6822 6823 } 6824 6825 enum NSEventModifierFlags : NSUInteger { 6826 NSEventModifierFlagCapsLock = 1 << 16, 6827 NSEventModifierFlagShift = 1 << 17, 6828 NSEventModifierFlagControl = 1 << 18, 6829 NSEventModifierFlagOption = 1 << 19, // aka Alt 6830 NSEventModifierFlagCommand = 1 << 20, // aka super 6831 NSEventModifierFlagNumericPad = 1 << 21, 6832 NSEventModifierFlagHelp = 1 << 22, 6833 NSEventModifierFlagFunction = 1 << 23, 6834 NSEventModifierFlagDeviceIndependentFlagsMask = 0xffff0000UL 6835 } 6836 6837 extern class NSEvent : NSObject { 6838 NSEventType type() @selector("type"); 6839 6840 NSPoint locationInWindow() @selector("locationInWindow"); 6841 NSTimeInterval timestamp() @selector("timestamp"); 6842 NSWindow window() @selector("window"); // note: nullable 6843 NSEventModifierFlags modifierFlags() @selector("modifierFlags"); 6844 6845 NSString characters() @selector("characters"); 6846 NSString charactersIgnoringModifiers() @selector("charactersIgnoringModifiers"); 6847 ushort keyCode() @selector("keyCode"); 6848 ushort specialKey() @selector("specialKey"); 6849 6850 static NSUInteger pressedMouseButtons() @selector("pressedMouseButtons"); 6851 NSPoint locationInWindow() @selector("locationInWindow"); // in screen coordinates 6852 static NSPoint mouseLocation() @selector("mouseLocation"); // in screen coordinates 6853 NSInteger buttonNumber() @selector("buttonNumber"); 6854 6855 CGFloat deltaX() @selector("deltaX"); 6856 CGFloat deltaY() @selector("deltaY"); 6857 CGFloat deltaZ() @selector("deltaZ"); 6858 6859 bool hasPreciseScrollingDeltas() @selector("hasPreciseScrollingDeltas"); 6860 6861 CGFloat scrollingDeltaX() @selector("scrollingDeltaX"); 6862 CGFloat scrollingDeltaY() @selector("scrollingDeltaY"); 6863 6864 // @property(getter=isDirectionInvertedFromDevice, readonly) BOOL directionInvertedFromDevice; 6865 } 6866 6867 extern /* final */ class NSTimer : NSObject { // the docs say don't subclass this, but making it final breaks the bridge 6868 override static NSTimer alloc() @selector("alloc"); 6869 override NSTimer init() @selector("init"); 6870 6871 static NSTimer schedule(NSTimeInterval timeIntervalInSeconds, NSid target, SEL selector, NSid userInfo, bool repeats) @selector("scheduledTimerWithTimeInterval:target:selector:userInfo:repeats:"); 6872 6873 void fire() @selector("fire"); 6874 void invalidate() @selector("invalidate"); 6875 6876 bool valid() @selector("isValid"); 6877 // @property(copy) NSDate *fireDate; 6878 NSTimeInterval timeInterval() @selector("timeInterval"); 6879 NSid userInfo() @selector("userInfo"); 6880 6881 NSTimeInterval tolerance() @selector("tolerance"); 6882 NSTimeInterval tolerance(NSTimeInterval) @selector("setTolerance:"); 6883 } 6884 6885 alias NSTimeInterval = double; 6886 6887 extern class NSResponder : NSObject { 6888 NSMenu menu() @selector("menu"); 6889 void menu(NSMenu menu) @selector("setMenu:"); 6890 6891 void keyDown(NSEvent event) @selector("keyDown:"); 6892 void keyUp(NSEvent event) @selector("keyUp:"); 6893 6894 // - (void)interpretKeyEvents:(NSArray<NSEvent *> *)eventArray; 6895 6896 void mouseDown(NSEvent event) @selector("mouseDown:"); 6897 void mouseDragged(NSEvent event) @selector("mouseDragged:"); 6898 void mouseUp(NSEvent event) @selector("mouseUp:"); 6899 void mouseMoved(NSEvent event) @selector("mouseMoved:"); 6900 void mouseEntered(NSEvent event) @selector("mouseEntered:"); 6901 void mouseExited(NSEvent event) @selector("mouseExited:"); 6902 6903 void rightMouseDown(NSEvent event) @selector("rightMouseDown:"); 6904 void rightMouseDragged(NSEvent event) @selector("rightMouseDragged:"); 6905 void rightMouseUp(NSEvent event) @selector("rightMouseUp:"); 6906 6907 void otherMouseDown(NSEvent event) @selector("otherMouseDown:"); 6908 void otherMouseDragged(NSEvent event) @selector("otherMouseDragged:"); 6909 void otherMouseUp(NSEvent event) @selector("otherMouseUp:"); 6910 6911 void scrollWheel(NSEvent event) @selector("scrollWheel:"); 6912 6913 // touch events should also be here btw among others 6914 } 6915 6916 extern class NSApplication : NSResponder { 6917 static NSApplication shared_() @selector("sharedApplication"); 6918 6919 NSApplicationDelegate delegate_() @selector("delegate"); 6920 void delegate_(NSApplicationDelegate) @selector("setDelegate:"); 6921 6922 bool setActivationPolicy(NSApplicationActivationPolicy activationPolicy) @selector("setActivationPolicy:"); 6923 6924 void activateIgnoringOtherApps(bool flag) @selector("activateIgnoringOtherApps:"); 6925 6926 @property NSMenu mainMenu() @selector("mainMenu"); 6927 @property NSMenu mainMenu(NSMenu) @selector("setMainMenu:"); 6928 6929 void run() @selector("run"); 6930 6931 void terminate(void*) @selector("terminate:"); 6932 } 6933 6934 extern interface NSApplicationDelegate { 6935 void applicationWillFinishLaunching(NSNotification notification) @selector("applicationWillFinishLaunching:"); 6936 void applicationDidFinishLaunching(NSNotification notification) @selector("applicationDidFinishLaunching:"); 6937 bool applicationShouldTerminateAfterLastWindowClosed(NSNotification notification) @selector("applicationShouldTerminateAfterLastWindowClosed:"); 6938 } 6939 6940 extern class NSNotification : NSObject { 6941 @property NSid object() @selector("object"); 6942 } 6943 6944 enum NSApplicationActivationPolicy : ptrdiff_t { 6945 /* The application is an ordinary app that appears in the Dock and may have a user interface. This is the default for bundled apps, unless overridden in the Info.plist. */ 6946 regular, 6947 6948 /* The application does not appear in the Dock and does not have a menu bar, but it may be activated programmatically or by clicking on one of its windows. This corresponds to LSUIElement=1 in the Info.plist. */ 6949 accessory, 6950 6951 /* The application does not appear in the Dock and may not create windows or be activated. This corresponds to LSBackgroundOnly=1 in the Info.plist. This is also the default for unbundled executables that do not have Info.plists. */ 6952 prohibited 6953 } 6954 6955 extern class NSGraphicsContext : NSObject { 6956 static NSGraphicsContext currentContext() @selector("currentContext"); 6957 NSGraphicsContext graphicsPort() @selector("graphicsPort"); 6958 } 6959 6960 extern class NSMenu : NSObject { 6961 override static NSMenu alloc() @selector("alloc"); 6962 6963 override NSMenu init() @selector("init"); 6964 NSMenu init(NSString title) @selector("initWithTitle:"); 6965 6966 void setSubmenu(NSMenu menu, NSMenuItem item) @selector("setSubmenu:forItem:"); 6967 void addItem(NSMenuItem newItem) @selector("addItem:"); 6968 6969 NSMenuItem addItem( 6970 NSString title, 6971 SEL selector, 6972 NSString charCode 6973 ) @selector("addItemWithTitle:action:keyEquivalent:"); 6974 } 6975 6976 extern class NSMenuItem : NSObject { 6977 override static NSMenuItem alloc() @selector("alloc"); 6978 override NSMenuItem init() @selector("init"); 6979 6980 NSMenuItem init( 6981 NSString title, 6982 SEL selector, 6983 NSString charCode 6984 ) @selector("initWithTitle:action:keyEquivalent:"); 6985 6986 void enabled(bool) @selector("setEnabled:"); 6987 6988 NSResponder target(NSResponder) @selector("setTarget:"); 6989 } 6990 6991 enum NSWindowStyleMask : size_t { 6992 borderless = 0, 6993 titled = 1 << 0, 6994 closable = 1 << 1, 6995 miniaturizable = 1 << 2, 6996 resizable = 1 << 3, 6997 6998 /* Specifies a window with textured background. Textured windows generally don't draw a top border line under the titlebar/toolbar. To get that line, use the NSUnifiedTitleAndToolbarWindowMask mask. 6999 */ 7000 texturedBackground = 1 << 8, 7001 7002 /* Specifies a window whose titlebar and toolbar have a unified look - that is, a continuous background. Under the titlebar and toolbar a horizontal separator line will appear. 7003 */ 7004 unifiedTitleAndToolbar = 1 << 12, 7005 7006 /* When set, the window will appear full screen. This mask is automatically toggled when toggleFullScreen: is called. 7007 */ 7008 fullScreen = 1 << 14, 7009 7010 /* If set, the contentView will consume the full size of the window; it can be combined with other window style masks, but is only respected for windows with a titlebar. 7011 Utilizing this mask opts-in to layer-backing. Utilize the contentLayoutRect or auto-layout contentLayoutGuide to layout views underneath the titlebar/toolbar area. 7012 */ 7013 fullSizeContentView = 1 << 15, 7014 7015 /* The following are only applicable for NSPanel (or a subclass thereof) 7016 */ 7017 utilityWindow = 1 << 4, 7018 docModalWindow = 1 << 6, 7019 nonactivatingPanel = 1 << 7, // Specifies that a panel that does not activate the owning application 7020 hUDWindow = 1 << 13 // Specifies a heads up display panel 7021 } 7022 7023 extern class NSWindow : NSObject { 7024 override static NSWindow alloc() @selector("alloc"); 7025 7026 override NSWindow init() @selector("init"); 7027 7028 NSWindow initWithContentRect( 7029 NSRect contentRect, 7030 NSWindowStyleMask style, 7031 NSBackingStoreType bufferingType, 7032 bool flag 7033 ) @selector("initWithContentRect:styleMask:backing:defer:"); 7034 7035 void makeKeyAndOrderFront(NSid sender) @selector("makeKeyAndOrderFront:"); 7036 NSView contentView() @selector("contentView"); 7037 void contentView(NSView view) @selector("setContentView:"); 7038 void orderFrontRegardless() @selector("orderFrontRegardless"); 7039 void center() @selector("center"); 7040 7041 NSRect frame() @selector("frame"); 7042 7043 NSRect contentRectForFrameRect(NSRect frameRect) @selector("contentRectForFrameRect:"); 7044 7045 NSString title() @selector("title"); 7046 void title(NSString value) @selector("setTitle:"); 7047 7048 void close() @selector("close"); 7049 7050 NSWindowDelegate delegate_() @selector("delegate"); 7051 void delegate_(NSWindowDelegate) @selector("setDelegate:"); 7052 7053 void setBackgroundColor(NSColor color) @selector("setBackgroundColor:"); 7054 } 7055 7056 extern interface NSWindowDelegate { 7057 @optional: 7058 void windowDidResize(NSNotification notification) @selector("windowDidResize:"); 7059 7060 NSSize windowWillResize(NSWindow sender, NSSize frameSize) @selector("windowWillResize:toSize:"); 7061 7062 void windowWillClose(NSNotification notification) @selector("windowWillClose:"); 7063 } 7064 7065 extern class NSView : NSResponder { 7066 override NSView init() @selector("init"); 7067 NSView initWithFrame(NSRect frameRect) @selector("initWithFrame:"); 7068 7069 void addSubview(NSView view) @selector("addSubview:"); 7070 7071 bool wantsLayer() @selector("wantsLayer"); 7072 void wantsLayer(bool value) @selector("setWantsLayer:"); 7073 7074 CALayer layer() @selector("layer"); 7075 void uiDelegate(NSObject) @selector("setUIDelegate:"); 7076 7077 void drawRect(NSRect rect) @selector("drawRect:"); 7078 bool isFlipped() @selector("isFlipped"); 7079 bool acceptsFirstResponder() @selector("acceptsFirstResponder"); 7080 bool setNeedsDisplay(bool) @selector("setNeedsDisplay:"); 7081 7082 // DO NOT USE: https://issues.dlang.org/show_bug.cgi?id=19017 7083 // an asm { pop RAX; } after getting the struct can kinda hack around this but still 7084 @property NSRect frame() @selector("frame"); 7085 @property NSRect frame(NSRect rect) @selector("setFrame:"); 7086 7087 void setFrameSize(NSSize newSize) @selector("setFrameSize:"); 7088 void setFrameOrigin(NSPoint newOrigin) @selector("setFrameOrigin:"); 7089 7090 void addSubview(NSView what) @selector("addSubview:"); 7091 void removeFromSuperview() @selector("removeFromSuperview"); 7092 } 7093 7094 extern class NSFont : NSObject { 7095 void set() @selector("set"); // sets it into the current graphics context 7096 void setInContext(NSGraphicsContext context) @selector("setInContext:"); 7097 7098 static NSFont fontWithName(NSString fontName, CGFloat fontSize) @selector("fontWithName:size:"); 7099 // fontWithDescriptor too 7100 // fontWithName and matrix too 7101 static NSFont systemFontOfSize(CGFloat fontSize) @selector("systemFontOfSize:"); 7102 // among others 7103 7104 @property CGFloat pointSize() @selector("pointSize"); 7105 @property bool isFixedPitch() @selector("isFixedPitch"); 7106 // fontDescriptor 7107 @property NSString displayName() @selector("displayName"); 7108 7109 @property CGFloat ascender() @selector("ascender"); 7110 @property CGFloat descender() @selector("descender"); // note it is negative 7111 @property CGFloat capHeight() @selector("capHeight"); 7112 @property CGFloat leading() @selector("leading"); 7113 @property CGFloat xHeight() @selector("xHeight"); 7114 // among many more 7115 } 7116 7117 extern class NSColor : NSObject { 7118 override static NSColor alloc() @selector("alloc"); 7119 static NSColor redColor() @selector("redColor"); 7120 static NSColor whiteColor() @selector("whiteColor"); 7121 7122 CGColorRef CGColor() @selector("CGColor"); 7123 } 7124 7125 extern class CALayer : NSObject { 7126 CGFloat borderWidth() @selector("borderWidth"); 7127 void borderWidth(CGFloat value) @selector("setBorderWidth:"); 7128 7129 CGColorRef borderColor() @selector("borderColor"); 7130 void borderColor(CGColorRef) @selector("setBorderColor:"); 7131 } 7132 7133 7134 extern class NSViewController : NSObject { 7135 NSView view() @selector("view"); 7136 void view(NSView view) @selector("setView:"); 7137 } 7138 7139 enum NSBackingStoreType : size_t { 7140 retained = 0, 7141 nonretained = 1, 7142 buffered = 2 7143 } 7144 7145 enum NSStringEncoding : NSUInteger { 7146 NSASCIIStringEncoding = 1, /* 0..127 only */ 7147 NSUTF8StringEncoding = 4, 7148 NSUnicodeStringEncoding = 10, 7149 7150 NSUTF16StringEncoding = NSUnicodeStringEncoding, 7151 NSUTF16BigEndianStringEncoding = 0x90000100, 7152 NSUTF16LittleEndianStringEncoding = 0x94000100, 7153 NSUTF32StringEncoding = 0x8c000100, 7154 NSUTF32BigEndianStringEncoding = 0x98000100, 7155 NSUTF32LittleEndianStringEncoding = 0x9c000100 7156 } 7157 7158 7159 struct CGColor; 7160 alias CGColorRef = CGColor*; 7161 7162 // note on the watch os it is float, not double 7163 alias CGFloat = double; 7164 7165 struct NSPoint { 7166 CGFloat x; 7167 CGFloat y; 7168 } 7169 7170 struct NSSize { 7171 CGFloat width; 7172 CGFloat height; 7173 } 7174 7175 struct NSRect { 7176 NSPoint origin; 7177 NSSize size; 7178 } 7179 7180 alias NSPoint CGPoint; 7181 alias NSSize CGSize; 7182 alias NSRect CGRect; 7183 7184 pragma(inline, true) NSPoint NSMakePoint(CGFloat x, CGFloat y) { 7185 NSPoint p; 7186 p.x = x; 7187 p.y = y; 7188 return p; 7189 } 7190 7191 pragma(inline, true) NSSize NSMakeSize(CGFloat w, CGFloat h) { 7192 NSSize s; 7193 s.width = w; 7194 s.height = h; 7195 return s; 7196 } 7197 7198 pragma(inline, true) NSRect NSMakeRect(CGFloat x, CGFloat y, CGFloat w, CGFloat h) { 7199 NSRect r; 7200 r.origin.x = x; 7201 r.origin.y = y; 7202 r.size.width = w; 7203 r.size.height = h; 7204 return r; 7205 } 7206 7207 7208 } 7209 7210 // helper raii refcount object 7211 static if(UseCocoa) 7212 struct MacString { 7213 union { 7214 // must be wrapped cuz of bug in dmd 7215 // referencing an init symbol when it should 7216 // just be null. but the union makes it work 7217 NSString s; 7218 } 7219 7220 // FIXME: if a string literal it would be kinda nice to use 7221 // the other function. but meh 7222 7223 this(scope const char[] str) { 7224 this.s = NSString.alloc.initWithBytes( 7225 cast(const(ubyte)*) str.ptr, 7226 str.length, 7227 NSStringEncoding.NSUTF8StringEncoding 7228 ); 7229 } 7230 7231 NSString borrow() { 7232 return s; 7233 } 7234 7235 this(this) { 7236 if(s !is null) 7237 s.retain(); 7238 } 7239 7240 ~this() { 7241 if(s !is null) { 7242 s.release(); 7243 s = null; 7244 } 7245 } 7246 } 7247 7248 extern(C) void NSLog(NSString, ...); 7249 extern(C) SEL sel_registerName(const(char)* str); 7250 7251 extern (Objective-C) __gshared NSApplication NSApp_; 7252 7253 NSApplication NSApp() { 7254 if(NSApp_ is null) 7255 NSApp_ = NSApplication.shared_; 7256 return NSApp_; 7257 } 7258 7259 // hacks to work around compiler bug 7260 extern(C) __gshared void* _D4arsd4core17NSGraphicsContext7__ClassZ = null; 7261 extern(C) __gshared void* _D4arsd4core6NSView7__ClassZ = null; 7262 extern(C) __gshared void* _D4arsd4core8NSWindow7__ClassZ = null; 7263 }