diff --git a/Cargo.lock b/Cargo.lock index 476cf221e7..132f30d794 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -23,7 +23,7 @@ version = "0.1.0" dependencies = [ "aml", "amlserde", - "arrayvec 0.7.6", + "arrayvec", "common", "libredox", "log", @@ -32,7 +32,7 @@ dependencies = [ "parking_lot 0.12.3", "plain", "redox-daemon", - "redox-scheme", + "redox-scheme 0.4.0", "redox_event", "redox_syscall", "ron", @@ -69,7 +69,7 @@ dependencies = [ name = "alxd" version = "0.1.0" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", "common", "libredox", "redox-daemon", @@ -125,15 +125,9 @@ dependencies = [ [[package]] name = "anyhow" -version = "1.0.92" +version = "1.0.97" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "74f37166d7d48a0284b99dd824694c26119c700b53bf0d1540cdb147dbdaaf13" - -[[package]] -name = "arrayvec" -version = "0.5.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "23b62fc65de8e4e7f52534fb52b0f3ed04746ae267519eef2a83941e8085068b" +checksum = "dcfed56ad506cb2c684a14971b8861fdc3baaaae314b9e5f9bb532cbe3ba7a4f" [[package]] name = "arrayvec" @@ -193,7 +187,7 @@ dependencies = [ "orbclient", "pcid", "redox-daemon", - "redox-scheme", + "redox-scheme 0.4.0", "redox_syscall", ] @@ -220,9 +214,9 @@ checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a" [[package]] name = "bitflags" -version = "2.6.0" +version = "2.9.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b048fb63fd8b5923fc5aa7b340d8e156aec7ec02f0c78fa8a6ddc2613f6f71de" +checksum = "5c8214115b7bf84099f1309324e63141d4c5d7cc26862f97a0a857dbefe165bd" dependencies = [ "serde", ] @@ -241,9 +235,9 @@ dependencies = [ [[package]] name = "bumpalo" -version = "3.16.0" +version = "3.17.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "79296716171880943b8470b5f8d03aa55eb2e645a4874bdbb28adb49162e012c" +checksum = "1628fb46dfa0b37568d12e5edd512553eccf6a22a78e8bde00bb4aed84d5bdbf" [[package]] name = "byteorder" @@ -253,9 +247,9 @@ checksum = "1fd0f2584146f6f2ef48085050886acf353beff7305ebd1ae69500e27c67f64b" [[package]] name = "cc" -version = "1.1.34" +version = "1.2.17" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "67b9470d453346108f93a59222a9a1a5724db32d0a4727b7ab7ace4b4d822dc9" +checksum = "1fcb57c740ae1daf453ae85f16e37396f672b039e00d9d866e07ddb24e328e3a" dependencies = [ "shlex", ] @@ -284,16 +278,16 @@ dependencies = [ [[package]] name = "chrono" -version = "0.4.38" +version = "0.4.40" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a21f936df1771bf62b77f047b726c4625ff2e8aa607c01ec06e5a05bd8463401" +checksum = "1a7964611d71df112cb1730f2ee67324fcf4d0fc6606acbbe9bfe06df124637c" dependencies = [ "android-tzdata", "iana-time-zone", "js-sys", "num-traits", "wasm-bindgen", - "windows-targets", + "windows-link", ] [[package]] @@ -363,11 +357,11 @@ dependencies = [ [[package]] name = "crossbeam-queue" -version = "0.3.11" +version = "0.3.12" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "df0346b5d5e76ac2fe4e327c5fd1118d6be7c51dfb18f9b7922923f287471e35" +checksum = "0f58bbc28f91df819d0aa2a2c00cd19754769c2fad90579b3592b1c9ba7a3115" dependencies = [ - "crossbeam-utils 0.8.20", + "crossbeam-utils 0.8.21", ] [[package]] @@ -383,17 +377,20 @@ dependencies = [ [[package]] name = "crossbeam-utils" -version = "0.8.20" +version = "0.8.21" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "22ec99545bb0ed0ea7bb9b8e1e9122ea386ff8a48c0922e43f36d45ab09e0e80" +checksum = "d0a5c400df2834b80a4c3327b3aad3a4c4cd4de0629063962b03235697506a28" [[package]] name = "driver-block" version = "0.1.0" dependencies = [ + "executor", + "futures", "libredox", + "log", "partitionlib", - "redox-scheme", + "redox-scheme 0.5.0", "redox_syscall", ] @@ -406,7 +403,7 @@ dependencies = [ "inputd", "libredox", "log", - "redox-scheme", + "redox-scheme 0.4.0", "redox_syscall", ] @@ -415,7 +412,7 @@ name = "driver-network" version = "0.1.0" dependencies = [ "libredox", - "redox-scheme", + "redox-scheme 0.4.0", "redox_syscall", ] @@ -423,7 +420,7 @@ dependencies = [ name = "e1000d" version = "0.1.0" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", "common", "driver-network", "libredox", @@ -435,9 +432,18 @@ dependencies = [ [[package]] name = "equivalent" -version = "1.0.1" +version = "1.0.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5443807d6dff69373d433ab9ef5378ad8df50ca6298caf15de6e52e24aaf54d5" +checksum = "877a4ace8713b0bcf2a4e7eec82529c029f1d0619886d18145fea96c3ffe5c0f" + +[[package]] +name = "executor" +version = "0.1.0" +dependencies = [ + "log", + "redox_event", + "slab", +] [[package]] name = "fbbootlogd" @@ -450,7 +456,7 @@ dependencies = [ "orbclient", "ransid", "redox-daemon", - "redox-scheme", + "redox-scheme 0.4.0", "redox_event", "redox_syscall", ] @@ -466,7 +472,7 @@ dependencies = [ "orbclient", "ransid", "redox-daemon", - "redox-scheme", + "redox-scheme 0.4.0", "redox_event", "redox_syscall", ] @@ -550,7 +556,7 @@ checksum = "162ee34ebcb7c64a8abebc059ce0fee27c2262618d7b60ed8faf72fef13c3650" dependencies = [ "proc-macro2", "quote", - "syn 2.0.87", + "syn 2.0.100", ] [[package]] @@ -585,12 +591,13 @@ dependencies = [ [[package]] name = "getrandom" -version = "0.2.15" +version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c4567c8db10ae91089c99af84c68c38da3ec2f087c3f82960bcdbf3656b6f4d7" +checksum = "73fea8450eea4bac3940448fb7ae50d91f034f941199fcd9d909a5a07aa455f0" dependencies = [ "cfg-if 1.0.0", "libc", + "r-efi", "wasi", ] @@ -600,7 +607,7 @@ version = "3.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8283e7331b8c93b9756e0cfdbcfb90312852f953c6faf9bf741e684cc3b6ad69" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", "crc", "log", "uuid", @@ -617,9 +624,9 @@ dependencies = [ [[package]] name = "hashbrown" -version = "0.15.1" +version = "0.15.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3a9bfc1af68b1726ea47d3d5109de126281def866b33970e10fbab11b5dafab3" +checksum = "bf151400ff0baff5465007dd2f3e717f3fe502074ca563069ce3a6629d07b289" [[package]] name = "hermit-abi" @@ -650,14 +657,15 @@ dependencies = [ [[package]] name = "iana-time-zone" -version = "0.1.61" +version = "0.1.62" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "235e081f3925a06703c2d0117ea8b91f042756fd6e7a6e5d901e8ca1a996b220" +checksum = "b2fd658b06e56721792c5df4475705b6cda790e9298d19d2f8af083457bcd127" dependencies = [ "android_system_properties", "core-foundation-sys", "iana-time-zone-haiku", "js-sys", + "log", "wasm-bindgen", "windows-core", ] @@ -689,7 +697,7 @@ dependencies = [ name = "ihdad" version = "0.1.0" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", "common", "libredox", "log", @@ -702,9 +710,9 @@ dependencies = [ [[package]] name = "indexmap" -version = "2.6.0" +version = "2.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "707907fe3c25f5424cce2cb7e1cbcafee6bdbe735ca90ef77c29e84591e5b9da" +checksum = "3954d50fe15b02142bf25d3b8bdadb634ec3948f103d04ffe3031bc8fe9d7058" dependencies = [ "equivalent", "hashbrown", @@ -720,21 +728,21 @@ dependencies = [ "log", "orbclient", "redox-daemon", - "redox-scheme", + "redox-scheme 0.4.0", "redox_syscall", ] [[package]] name = "itoa" -version = "1.0.11" +version = "1.0.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "49f1f14873335454500d59611f1cf4a4b0f786f9ac11f4312a78e4cf2566695b" +checksum = "4a5f13b858c8d314ee3e8f639011f7ccefe71f97f96e50151fb991f267928e2c" [[package]] name = "ixgbed" version = "1.0.0" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", "common", "driver-network", "libredox", @@ -746,10 +754,11 @@ dependencies = [ [[package]] name = "js-sys" -version = "0.3.72" +version = "0.3.77" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6a88f1bda2bd75b0452a14784937d796722fdebfe50df998aeb3f0b7603019a9" +checksum = "1cfaf33c695fc6e08064efbc1f72ec937429614f25eef83af942d0e227c3a28f" dependencies = [ + "once_cell", "wasm-bindgen", ] @@ -761,9 +770,9 @@ checksum = "bbd2bcb4c963f2ddae06a2efc7e9f3591312473c50c6685e1f298068316e66fe" [[package]] name = "libc" -version = "0.2.161" +version = "0.2.171" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8e9489c2807c139ffd9c1794f4af0ebe86a828db53ecdc7fea2111d0fed085d1" +checksum = "c19937216e9d3aa9956d9bb8dfc0b0c8beb6058fc4f7a4dc4d850edf86a237d6" [[package]] name = "libredox" @@ -771,7 +780,7 @@ version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c0ff37bd590ca25063e35af745c343cb7a0271906fb7b37e4813e8f79f00268d" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", "libc", "redox_syscall", ] @@ -800,9 +809,9 @@ dependencies = [ [[package]] name = "log" -version = "0.4.22" +version = "0.4.26" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a7a70ba024b9dc04c27ea2f0c0548feb474ec5c54bba33a7f72f873a39d07b24" +checksum = "30bde2b3dc3671ae49d8e2e9f044c7c005836e7a023ee57cffa25ab82764bb9e" [[package]] name = "maybe-uninit" @@ -846,26 +855,27 @@ checksum = "6aa2c4e539b869820a2b82e1aef6ff40aa85e65decdd5185e83fb4b1249cd00f" name = "nvmed" version = "0.1.0" dependencies = [ - "arrayvec 0.5.2", - "bitflags 1.3.2", + "arrayvec", + "bitflags 2.9.0", "common", - "crossbeam-channel", "driver-block", - "futures", + "executor", "libredox", "log", + "parking_lot 0.12.3", + "partitionlib", "pcid", "redox-daemon", "redox_event", "redox_syscall", - "smallvec 1.13.2", + "smallvec 1.14.0", ] [[package]] name = "once_cell" -version = "1.20.2" +version = "1.21.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1261fe7e33c73b354eab43b1273a57c8f967d0391e80353e51f764ac02cf6775" +checksum = "d75b0bedcc4fe52caa0e03d9f1151a323e4aa5e2d78ba3580400cd3c9e2bc4bc" [[package]] name = "orbclient" @@ -928,7 +938,7 @@ dependencies = [ "cfg-if 1.0.0", "libc", "redox_syscall", - "smallvec 1.13.2", + "smallvec 1.14.0", "windows-targets", ] @@ -948,7 +958,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c4325c6aa3cca3373503b1527e75756f9fbfe5fd76be4b4c8a143ee47430b8e0" dependencies = [ "bit_field", - "bitflags 2.6.0", + "bitflags 2.9.0", ] [[package]] @@ -965,7 +975,7 @@ dependencies = [ "pico-args", "plain", "redox-daemon", - "redox-scheme", + "redox-scheme 0.4.0", "redox_syscall", "serde", ] @@ -992,9 +1002,9 @@ checksum = "5be167a7af36ee22fe3115051bc51f6e6c7054c9348e28deb4f49bd6f705a315" [[package]] name = "pin-project-lite" -version = "0.2.15" +version = "0.2.16" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "915a1e146535de9163f3987b8944ed8cf49a18bb0056bcebcdcece385cece4ff" +checksum = "3b3cff922bd51709b605d9ead9aa71031d81447142d828eb4a6eba76fe619f9b" [[package]] name = "pin-utils" @@ -1010,9 +1020,9 @@ checksum = "b4596b6d070b27117e987119b4dac604f3c58cfb0b191112e24771b2faeac1a6" [[package]] name = "proc-macro2" -version = "1.0.89" +version = "1.0.94" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f139b0662de085916d1fb67d2b4169d1addddda1919e696f3252b740b629986e" +checksum = "a31971752e70b8b2686d7e46ec17fb38dad4051d94024c88df49b667caea9c84" dependencies = [ "unicode-ident", ] @@ -1034,13 +1044,19 @@ dependencies = [ [[package]] name = "quote" -version = "1.0.37" +version = "1.0.40" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b5b9d34b8991d19d98081b46eacdd8eb58c6f2b201139f7c5f643cc155a633af" +checksum = "1885c039570dc00dcb4ff087a89e185fd56bae234ddc7f056a945bf36467248d" dependencies = [ "proc-macro2", ] +[[package]] +name = "r-efi" +version = "5.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "74765f6d916ee2faa39bc8e68e4f3ed8949b48cccdac59983d287a7cb71ce9c5" + [[package]] name = "radium" version = "0.7.0" @@ -1111,7 +1127,7 @@ checksum = "81460b1526438123d16f0c968dbe42ba7f61e99645109b70e57864a8b66710fb" dependencies = [ "chrono", "log", - "smallvec 1.13.2", + "smallvec 1.14.0", "termion", ] @@ -1125,23 +1141,33 @@ dependencies = [ "redox_syscall", ] +[[package]] +name = "redox-scheme" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a96b9cfb034251dfb0aaa66a67059a7f0ea344039904d1d70cd36266af9c8a2f" +dependencies = [ + "libredox", + "redox_syscall", +] + [[package]] name = "redox_event" version = "0.4.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "69609faa5d5992247a4ef379917bb3e39be281405d6a0ccd4f942429400b956f" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", "libredox", ] [[package]] name = "redox_syscall" -version = "0.5.9" +version = "0.5.10" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "82b568323e98e49e2a0899dcee453dd679fae22d69adf9b11dd508d1549b7e2f" +checksum = "0b8c0c260b63a8219631167be35e6a988e9554dbd323f8bd08439c8ed1302bd1" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", ] [[package]] @@ -1164,9 +1190,9 @@ dependencies = [ [[package]] name = "regex-automata" -version = "0.4.8" +version = "0.4.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "368758f23274712b504848e9d5a6f010445cc8b87a7cdb4d7cbee666c1288da3" +checksum = "809e8dc61f6de73b46c85f4c96486310fe304c434cfa43669d7b40f711150908" dependencies = [ "aho-corasick", "memchr", @@ -1195,7 +1221,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b91f7eff05f748767f183df4320a63d6936e9c6107d97c9e6bdd9784f4289c94" dependencies = [ "base64 0.21.7", - "bitflags 2.6.0", + "bitflags 2.9.0", "serde", "serde_derive", ] @@ -1204,7 +1230,7 @@ dependencies = [ name = "rtl8139d" version = "0.1.0" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", "common", "driver-network", "libredox", @@ -1219,7 +1245,7 @@ dependencies = [ name = "rtl8168d" version = "0.1.0" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", "common", "driver-network", "libredox", @@ -1237,16 +1263,22 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "08d43f7aa6b08d49f382cde6a7982047c3426db949b1424bc4b7ec9ae12c6ce2" [[package]] -name = "ryu" -version = "1.0.18" +name = "rustversion" +version = "1.0.20" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f3cb5ba0dc43242ce17de99c180e96db90b235b8a9fdc9543c96d2209116bd9f" +checksum = "eded382c5f5f786b989652c49544c4877d9f015cc22e145a5ea8ea66c2921cd2" + +[[package]] +name = "ryu" +version = "1.0.20" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "28d3b2b1366ec20994f1fd18c3c594f05c5dd4bc44d8bb0c1c632c8d6829481f" [[package]] name = "sb16d" version = "0.1.0" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", "common", "libredox", "log", @@ -1307,29 +1339,29 @@ dependencies = [ [[package]] name = "serde" -version = "1.0.214" +version = "1.0.219" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f55c3193aca71c12ad7890f1785d2b73e1b9f63a0bbc353c08ef26fe03fc56b5" +checksum = "5f0e2c6ed6606019b4e29e69dbaba95b11854410e5347d525002456dbbb786b6" dependencies = [ "serde_derive", ] [[package]] name = "serde_derive" -version = "1.0.214" +version = "1.0.219" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "de523f781f095e28fa605cdce0f8307e451cc0fd14e2eb4cd2e98a355b147766" +checksum = "5b0276cf7f2c73365f7157c8123c21cd9a50fbbd844757af28ca1f5925fc2a00" dependencies = [ "proc-macro2", "quote", - "syn 2.0.87", + "syn 2.0.100", ] [[package]] name = "serde_json" -version = "1.0.132" +version = "1.0.140" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d726bfaff4b320266d395898905d0eba0345aae23b54aee3a737e260fd46db03" +checksum = "20068b6e96dc6c9bd23e01df8827e6c7e1f2fddd43c21810382803c136b99373" dependencies = [ "itoa", "memchr", @@ -1372,9 +1404,9 @@ dependencies = [ [[package]] name = "smallvec" -version = "1.13.2" +version = "1.14.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3c5e1a9a646d36c3599cd173a41282daf47c44583ad367b8e6837255952e5c67" +checksum = "7fcf8323ef1faaee30a44a340193b1ac6814fd9b7b4e88e9d4519a3e4abe1cfd" dependencies = [ "serde", ] @@ -1428,9 +1460,9 @@ dependencies = [ [[package]] name = "syn" -version = "2.0.87" +version = "2.0.100" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "25aa4ce346d03a6dcd68dd8b4010bcb74e54e62c90c573f394c46eae99aba32d" +checksum = "b09a44accad81e1ba1cd74a32461ba89dee89095ba17b32f5d03683b1b1fc2a0" dependencies = [ "proc-macro2", "quote", @@ -1445,9 +1477,9 @@ checksum = "55937e1799185b12863d447f42597ed69d9928686b8d88a1df17376a097d8369" [[package]] name = "termion" -version = "4.0.3" +version = "4.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7eaa98560e51a2cf4f0bb884d8b2098a9ea11ecf3b7078e9c68242c74cc923a7" +checksum = "6f359c854fbecc1ea65bc3683f1dcb2dce78b174a1ca7fda37acd1fff81df6ff" dependencies = [ "libc", "libredox", @@ -1466,22 +1498,22 @@ dependencies = [ [[package]] name = "thiserror" -version = "1.0.68" +version = "1.0.69" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "02dd99dc800bbb97186339685293e1cc5d9df1f8fae2d0aecd9ff1c77efea892" +checksum = "b6aaf5339b578ea85b50e080feb250a3e8ae8cfcdff9a461c9ec2904bc923f52" dependencies = [ "thiserror-impl", ] [[package]] name = "thiserror-impl" -version = "1.0.68" +version = "1.0.69" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a7c61ec9a6f64d2793d8a45faba21efbe3ced62a886d44c36a009b2b519b4c7e" +checksum = "4fee6c4efc90059e10f81e6d42c60a18f76588c3d74cb83a0b242a2b6c7504c1" dependencies = [ "proc-macro2", "quote", - "syn 2.0.87", + "syn 2.0.100", ] [[package]] @@ -1529,9 +1561,9 @@ dependencies = [ [[package]] name = "unicode-ident" -version = "1.0.13" +version = "1.0.18" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e91b56cd4cadaeb79bbf1a5645f6b4f8dc5bde8834ad5894a8db35fda9efa1fe" +checksum = "5a5f39404a5da50712a4c1eecf25e90dd62b613502b7e925fd4e4d19b5c96512" [[package]] name = "unicode-width" @@ -1551,7 +1583,7 @@ dependencies = [ name = "usbhidd" version = "0.1.0" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", "common", "inputd", "log", @@ -1594,9 +1626,9 @@ checksum = "8772a4ccbb4e89959023bc5b7cb8623a795caa7092d99f3aa9501b9484d4557d" [[package]] name = "uuid" -version = "1.11.0" +version = "1.16.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f8c5f0a0af699448548ad1a2fbf920fb4bee257eae39953ba95cb84891a0446a" +checksum = "458f7a779bf54acc9f347480ac654f68407d3aab21269a6e3c9f922acd9e2da9" dependencies = [ "getrandom", ] @@ -1666,7 +1698,7 @@ dependencies = [ name = "virtio-core" version = "0.1.0" dependencies = [ - "bitflags 2.6.0", + "bitflags 2.9.0", "common", "crossbeam-queue", "futures", @@ -1728,41 +1760,44 @@ dependencies = [ [[package]] name = "wasi" -version = "0.11.0+wasi-snapshot-preview1" +version = "0.14.2+wasi-0.2.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9c8d87e72b64a3b4db28d11ce29237c246188f4f51057d65a7eab63b7987e423" +checksum = "9683f9a5a998d873c0d21fcbe3c083009670149a8fab228644b8bd36b2c48cb3" +dependencies = [ + "wit-bindgen-rt", +] [[package]] name = "wasm-bindgen" -version = "0.2.95" +version = "0.2.100" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "128d1e363af62632b8eb57219c8fd7877144af57558fb2ef0368d0087bddeb2e" +checksum = "1edc8929d7499fc4e8f0be2262a241556cfc54a0bea223790e71446f2aab1ef5" dependencies = [ "cfg-if 1.0.0", "once_cell", + "rustversion", "wasm-bindgen-macro", ] [[package]] name = "wasm-bindgen-backend" -version = "0.2.95" +version = "0.2.100" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cb6dd4d3ca0ddffd1dd1c9c04f94b868c37ff5fac97c30b97cff2d74fce3a358" +checksum = "2f0a0651a5c2bc21487bde11ee802ccaf4c51935d0d3d42a6101f98161700bc6" dependencies = [ "bumpalo", "log", - "once_cell", "proc-macro2", "quote", - "syn 2.0.87", + "syn 2.0.100", "wasm-bindgen-shared", ] [[package]] name = "wasm-bindgen-macro" -version = "0.2.95" +version = "0.2.100" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e79384be7f8f5a9dd5d7167216f022090cf1f9ec128e6e6a482a2cb5c5422c56" +checksum = "7fe63fc6d09ed3792bd0897b314f53de8e16568c2b3f7982f468c0bf9bd0b407" dependencies = [ "quote", "wasm-bindgen-macro-support", @@ -1770,22 +1805,25 @@ dependencies = [ [[package]] name = "wasm-bindgen-macro-support" -version = "0.2.95" +version = "0.2.100" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "26c6ab57572f7a24a4985830b120de1594465e5d500f24afe89e16b4e833ef68" +checksum = "8ae87ea40c9f689fc23f209965b6fb8a99ad69aeeb0231408be24920604395de" dependencies = [ "proc-macro2", "quote", - "syn 2.0.87", + "syn 2.0.100", "wasm-bindgen-backend", "wasm-bindgen-shared", ] [[package]] name = "wasm-bindgen-shared" -version = "0.2.95" +version = "0.2.100" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "65fc09f10666a9f147042251e0dda9c18f166ff7de300607007e96bdebc1068d" +checksum = "1a05d73b933a847d6cccdda8f838a22ff101ad9bf93e33684f39c1f5f0eece3d" +dependencies = [ + "unicode-ident", +] [[package]] name = "winapi" @@ -1818,6 +1856,12 @@ dependencies = [ "windows-targets", ] +[[package]] +name = "windows-link" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "76840935b766e1b0a05c0066835fb9ec80071d4c09a16f6bd5f7e655e3c14c38" + [[package]] name = "windows-targets" version = "0.52.6" @@ -1891,6 +1935,15 @@ dependencies = [ "memchr", ] +[[package]] +name = "wit-bindgen-rt" +version = "0.39.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6f42320e61fe2cfd34354ecb597f86f413484a798ba44a8ca1165c58d42da6c1" +dependencies = [ + "bitflags 2.9.0", +] + [[package]] name = "wyz" version = "0.5.1" @@ -1915,13 +1968,13 @@ dependencies = [ "pcid", "plain", "redox-daemon", - "redox-scheme", + "redox-scheme 0.4.0", "redox_event", "redox_syscall", "regex", "serde", "serde_json", - "smallvec 1.13.2", + "smallvec 1.14.0", "thiserror", "toml 0.5.11", ] diff --git a/Cargo.toml b/Cargo.toml index 14dd17de9d..d27f9e6e46 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,7 +1,9 @@ [workspace] members = [ - "acpid", "common", + "executor", + + "acpid", "hwd", "pcid", "pcid-spawner", diff --git a/executor/Cargo.toml b/executor/Cargo.toml new file mode 100644 index 0000000000..0d5b12f070 --- /dev/null +++ b/executor/Cargo.toml @@ -0,0 +1,12 @@ +[package] +name = "executor" +authors = ["4lDO2 <4lDO2@protonmail.com>"] +version = "0.1.0" +edition = "2021" +license = "MIT" +description = "Async framework for queue-based HW interfaces" + +[dependencies] +log = "0.4" +redox_event = "0.4.1" +slab = "0.4.9" diff --git a/executor/src/lib.rs b/executor/src/lib.rs new file mode 100644 index 0000000000..6b54e07bc1 --- /dev/null +++ b/executor/src/lib.rs @@ -0,0 +1,396 @@ +use std::cell::{Cell, RefCell}; +use std::collections::{HashMap, VecDeque}; +use std::fmt::Debug; +use std::fs::File; +use std::future::{Future, IntoFuture}; +use std::hash::Hash; +use std::io::{Read, Write}; +use std::marker::PhantomData; +use std::os::fd::AsRawFd; +use std::panic::AssertUnwindSafe; +use std::pin::Pin; +use std::ptr::NonNull; +use std::rc::Rc; +use std::task; + +use event::{EventFlags, RawEventQueue}; +use slab::Slab; + +type EventUserData = usize; + +type FutIdx = usize; + +pub trait Hardware: Sized { + type CmdId: Clone + Copy + Debug + Hash + Eq + PartialEq; + type CqId: Clone + Copy + Debug + Hash + Eq + PartialEq; + type SqId: Clone + Copy + Debug + Hash + Eq + PartialEq; + type Sqe: Debug + Clone + Copy; + type Cqe; + type Iv: Clone + Copy + Debug; + + type GlobalCtxt; + + // TODO: the kernel should also do this automatically before sending EOI messages to the IC + fn mask_vector(ctxt: &Self::GlobalCtxt, iv: Self::Iv); + fn unmask_vector(ctxt: &Self::GlobalCtxt, iv: Self::Iv); + + fn set_sqe_cmdid(sqe: &mut Self::Sqe, id: Self::CmdId); + fn get_cqe_cmdid(cqe: &Self::Cqe) -> Self::CmdId; + + // TODO: support multiple SQs per CQ or vice versa? + fn sq_cq(ctxt: &Self::GlobalCtxt, id: Self::CqId) -> Self::SqId; + + fn current() -> Rc>; + fn vtable() -> &'static task::RawWakerVTable; + + fn try_submit( + ctxt: &Self::GlobalCtxt, + sq_id: Self::SqId, + success: impl FnOnce(Self::CmdId) -> Self::Sqe, + fail: impl FnOnce(), + ) -> Option<(Self::CqId, Self::CmdId)>; + fn poll_cqes(ctxt: &Self::GlobalCtxt, handle: impl FnMut(Self::CqId, Self::Cqe)); +} + +/// Async executor, single IV, thread-per-core architecture +pub struct LocalExecutor { + global_ctxt: Hw::GlobalCtxt, + + queue: RawEventQueue, + vector: Hw::Iv, + irq_handle: File, + intx: bool, + + // TODO: One IV and SQ/CQ per core (where the admin queue can be managed by the main thread). + awaiting_submission: RefCell>>, + awaiting_completion: + RefCell>)>>>, + + external_event: RefCell)>>, + next_user_data: Cell, + + ready_futures: RefCell>, + futures: RefCell + 'static>>>>, + is_polling: Cell, +} + +impl LocalExecutor { + pub fn register_external_event( + &self, + fd: usize, + flags: event::EventFlags, + ) -> ExternalEventSource { + let user_data = self.next_user_data.get(); + self.next_user_data.set(user_data.checked_add(1).unwrap()); + + self.queue + .subscribe(fd, user_data, flags) + .expect("failed to subscribe to event"); + + ExternalEventSource { + flags: event::EventFlags::empty(), + user_data, + _not_send_or_unpin: PhantomData, + } + } + pub fn current() -> Rc { + Hw::current() + } + pub fn poll(&self) -> usize { + assert!(!self.is_polling.replace(true)); + + let mut finished = 0; + + for future_idx in self.ready_futures.borrow_mut().drain(..) { + let waker = waker::(future_idx); + + let mut futures = self.futures.borrow_mut(); + let res = match std::panic::catch_unwind(AssertUnwindSafe(|| { + futures[future_idx] + .as_mut() + .poll(&mut task::Context::from_waker(&waker)) + })) { + Ok(r) => r, + Err(_) => { + log::error!("Task panicked!"); + core::mem::forget(futures.remove(future_idx)); + continue; + } + }; + if res.is_ready() { + drop(futures.remove(future_idx)); + finished += 1; + } + } + self.is_polling.set(false); + + finished + } + pub fn spawn(&self, fut: impl IntoFuture + 'static) { + let idx = self + .futures + .borrow_mut() + .insert(Box::pin(fut.into_future())); + self.ready_futures.borrow_mut().push_back(idx); + } + pub fn block_on<'a, O: 'a>(&self, fut: impl IntoFuture + 'a) -> O { + let retval = Rc::new(RefCell::new(None)); + + let retval2 = Rc::clone(&retval); + let idx = self.futures.borrow_mut().insert({ + let t1: Pin + 'a>> = Box::pin(async move { + *retval2.borrow_mut() = Some(fut.await); + }); + // SAFETY: Apart from the lifetimes, the types are exactly the same. We also know + // block_on simply cannot return without having fully awaited and dropped the future, + // even if that future panics (cf. the catch_unwind invocation). + let t2: Pin + 'static>> = + unsafe { std::mem::transmute(t1) }; + + t2 + }); + + self.ready_futures.borrow_mut().push_front(idx); + + loop { + let finished = self.poll(); + if retval.borrow().is_some() { + break; + } + if finished == 0 { + self.react(); + } + } + + let o = retval.borrow_mut().take().unwrap(); + o + } + fn react(&self) { + let event = self.queue.next_event().expect("failed to get next event"); + + if event.user_data != 0 { + let Some((fut_idx, flags_ptr)) = + self.external_event.borrow_mut().remove(&event.user_data) + else { + // Spurious event + return; + }; + unsafe { + flags_ptr + .as_ptr() + .write(event::EventFlags::from_bits_retain(event.flags)); + } + self.ready_futures.borrow_mut().push_back(fut_idx); + return; + } + + if self.intx { + let mut buf = [0_u8; core::mem::size_of::()]; + if (&self.irq_handle).read(&mut buf).unwrap() != 0 { + (&self.irq_handle).write(&buf).unwrap(); + } + } + + // TODO: The kernel should probably do the masking (when using MSI/MSI-X at least), which + // should happen before EOI messages to the interrupt controller. + Hw::mask_vector(&self.global_ctxt, self.vector); + Hw::poll_cqes(&self.global_ctxt, |cq_id, cqe| { + if let Some((fut_idx, comp_ptr)) = self + .awaiting_completion + .borrow_mut() + .get_mut(&cq_id) + .and_then(|per_cmd| per_cmd.remove(&Hw::get_cqe_cmdid(&cqe))) + { + unsafe { + comp_ptr.as_ptr().write(Some(cqe)); + } + self.ready_futures.borrow_mut().push_back(fut_idx); + + if let Some(submitting) = self + .awaiting_submission + .borrow_mut() + .get_mut(&Hw::sq_cq(&self.global_ctxt, cq_id)) + .and_then(|q| q.pop_front()) + { + self.ready_futures.borrow_mut().push_back(submitting); + } + } + }); + Hw::unmask_vector(&self.global_ctxt, self.vector); + } + pub async fn submit(&self, sq_id: Hw::SqId, cmd: Hw::Sqe) -> Hw::Cqe { + CqeFuture:: { + state: State::::Submitting { sq_id, cmd }, + comp: None, + _not_send: PhantomData, + } + .await + } +} + +struct CqeFuture { + pub state: State, + pub comp: Option, + pub _not_send: PhantomData<*const ()>, +} +enum State { + Submitting { sq_id: Hw::SqId, cmd: Hw::Sqe }, + Completing { cq_id: Hw::CqId, cmd_id: Hw::CmdId }, +} + +fn current_executor_and_idx( + cx: &mut task::Context<'_>, +) -> (Rc>, FutIdx) { + let executor = LocalExecutor::current(); + + let idx = cx.waker().data() as FutIdx; + assert_eq!( + cx.waker().vtable() as *const _, + Hw::vtable(), + "incompatible executor for CqeFuture" + ); + + (executor, idx) +} + +impl Future for CqeFuture { + type Output = Hw::Cqe; + + fn poll(self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> task::Poll { + let this = unsafe { self.get_unchecked_mut() }; + + let (executor, idx) = current_executor_and_idx::(cx); + + match this.state { + State::Submitting { sq_id, mut cmd } => { + let mut awaiting = executor.awaiting_submission.borrow_mut(); + + if let Some((cq_id, cmd_id)) = Hw::try_submit( + &executor.global_ctxt, + sq_id, + |cmd_id| { + Hw::set_sqe_cmdid(&mut cmd, cmd_id); + log::trace!("About to submit {cmd:?}"); + cmd + }, + || { + awaiting.entry(sq_id).or_default().push_back(idx); + }, + ) { + executor + .awaiting_completion + .borrow_mut() + .entry(cq_id) + .or_default() + .insert(cmd_id, (idx, (&mut this.comp).into())); + this.state = State::Completing { cq_id, cmd_id }; + } + task::Poll::Pending + } + State::Completing { cq_id, cmd_id } => match this.comp.take() { + Some(comp) => { + log::trace!("ready!"); + task::Poll::Ready(comp) + } + + // Shouldn't technically be possible + None => { + log::trace!("spurious poll"); + executor + .awaiting_completion + .borrow_mut() + .entry(cq_id) + .or_default() + .insert(cmd_id, (idx, (&mut this.comp).into())); + task::Poll::Pending + } + }, + } + } +} + +unsafe fn vt_clone(idx: *const ()) -> task::RawWaker { + task::RawWaker::new(idx, Hw::vtable()) +} +unsafe fn vt_drop(_idx: *const ()) {} +unsafe fn vt_wake(idx: *const ()) { + Hw::current() + .ready_futures + .borrow_mut() + .push_back(idx as FutIdx); +} + +fn waker(idx: FutIdx) -> task::Waker { + unsafe { task::Waker::from_raw(task::RawWaker::new(idx as *const (), Hw::vtable())) } +} +pub const fn vtable() -> task::RawWakerVTable { + task::RawWakerVTable::new(vt_clone::, vt_wake::, vt_wake::, vt_drop) +} + +pub struct ExternalEventSource { + flags: event::EventFlags, + user_data: EventUserData, + _not_send_or_unpin: PhantomData<(*const (), fn() -> Hw)>, +} +pub struct Event { + flags: event::EventFlags, + _not_send: PhantomData<*const ()>, +} +impl Event { + pub fn flags(&self) -> event::EventFlags { + self.flags + } +} +impl ExternalEventSource { + fn poll_next(self: Pin<&mut Self>, cx: &mut task::Context) -> task::Poll> { + let this = unsafe { self.get_unchecked_mut() }; + + let flags = std::mem::take(&mut this.flags); + + if flags.is_empty() { + let (executor, idx) = current_executor_and_idx::(cx); + executor + .external_event + .borrow_mut() + .insert(this.user_data, (idx, (&mut this.flags).into())); + return task::Poll::Pending; + } + task::Poll::Ready(Some(Event { + flags, + _not_send: PhantomData, + })) + } + pub async fn next(mut self: Pin<&mut Self>) -> Option { + core::future::poll_fn(|cx| self.as_mut().poll_next(cx)).await + } +} +pub fn init_raw( + global_ctxt: Hw::GlobalCtxt, + vector: Hw::Iv, + intx: bool, + irq_handle: File, +) -> LocalExecutor { + let queue = RawEventQueue::new().expect("failed to allocate event queue for local executor"); + + // TODO: Multiple CPUs + queue + .subscribe(irq_handle.as_raw_fd() as usize, 0, EventFlags::READ) + .expect("failed to subscribe to IRQ event"); + + LocalExecutor { + global_ctxt, + + queue, + vector, + intx, + irq_handle, + + awaiting_submission: RefCell::new(HashMap::new()), + awaiting_completion: RefCell::new(HashMap::new()), + external_event: RefCell::new(HashMap::new()), + next_user_data: Cell::new(1), + ready_futures: RefCell::new(VecDeque::new()), + futures: RefCell::new(Slab::with_capacity(16)), + is_polling: Cell::new(false), + } +} diff --git a/storage/ahcid/src/ahci/disk_ata.rs b/storage/ahcid/src/ahci/disk_ata.rs index bb52a4571d..3494d8624d 100644 --- a/storage/ahcid/src/ahci/disk_ata.rs +++ b/storage/ahcid/src/ahci/disk_ata.rs @@ -144,6 +144,7 @@ impl DiskATA { self.request_opt = Some(request); + // TODO: support async internally return Ok(None); } else { // Done @@ -162,21 +163,21 @@ impl Disk for DiskATA { self.size } - fn read(&mut self, block: u64, buffer: &mut [u8]) -> Result> { + async fn read(&mut self, block: u64, buffer: &mut [u8]) -> Result { //TODO: FIGURE OUT WHY INTERRUPTS CAUSE HANGS loop { match self.request(block, BufferKind::Read(buffer))? { - Some(count) => return Ok(Some(count)), + Some(count) => return Ok(count), None => std::thread::yield_now(), } } } - fn write(&mut self, block: u64, buffer: &[u8]) -> Result> { + async fn write(&mut self, block: u64, buffer: &[u8]) -> Result { //TODO: FIGURE OUT WHY INTERRUPTS CAUSE HANGS loop { match self.request(block, BufferKind::Write(buffer))? { - Some(count) => return Ok(Some(count)), + Some(count) => return Ok(count), None => std::thread::yield_now(), } } diff --git a/storage/ahcid/src/ahci/disk_atapi.rs b/storage/ahcid/src/ahci/disk_atapi.rs index b92112c36e..9571ed9e0e 100644 --- a/storage/ahcid/src/ahci/disk_atapi.rs +++ b/storage/ahcid/src/ahci/disk_atapi.rs @@ -77,7 +77,7 @@ impl Disk for DiskATAPI { u64::from(self.blk_count) * u64::from(self.blk_size) } - fn read(&mut self, block: u64, buffer: &mut [u8]) -> Result> { + async fn read(&mut self, block: u64, buffer: &mut [u8]) -> Result { // TODO: Handle audio CDs, which use special READ CD command let blk_len = self.blk_size; @@ -139,10 +139,10 @@ impl Disk for DiskATAPI { sector += sectors - sector; } - Ok(Some((sector * blk_len) as usize)) + Ok((sector * blk_len) as usize) } - fn write(&mut self, _block: u64, _buffer: &[u8]) -> Result> { + async fn write(&mut self, _block: u64, _buffer: &[u8]) -> Result { Err(Error::new(EBADF)) // TODO: Implement writing } } diff --git a/storage/ahcid/src/ahci/mod.rs b/storage/ahcid/src/ahci/mod.rs index 73e66c5ec5..4d8cc8c040 100644 --- a/storage/ahcid/src/ahci/mod.rs +++ b/storage/ahcid/src/ahci/mod.rs @@ -11,27 +11,58 @@ pub mod disk_atapi; pub mod fis; pub mod hba; -pub fn disks(base: usize, name: &str) -> (&'static mut HbaMem, Vec>) { +pub enum AnyDisk { + Ata(DiskATA), + Atapi(DiskATAPI), +} +impl Disk for AnyDisk { + fn block_size(&self) -> u32 { + match self { + Self::Ata(a) => a.block_size(), + Self::Atapi(a) => a.block_size(), + } + } + fn size(&self) -> u64 { + match self { + Self::Ata(a) => a.size(), + Self::Atapi(a) => a.size(), + } + } + async fn read(&mut self, base: u64, buffer: &mut [u8]) -> syscall::Result { + match self { + Self::Ata(a) => a.read(base, buffer).await, + Self::Atapi(a) => a.read(base, buffer).await, + } + } + async fn write(&mut self, base: u64, buffer: &[u8]) -> syscall::Result { + match self { + Self::Ata(a) => a.write(base, buffer).await, + Self::Atapi(a) => a.write(base, buffer).await, + } + } +} + +pub fn disks(base: usize, name: &str) -> (&'static mut HbaMem, Vec) { let hba_mem = unsafe { &mut *(base as *mut HbaMem) }; hba_mem.init(); let pi = hba_mem.pi.read(); - let disks: Vec> = (0..hba_mem.ports.len()) + let disks: Vec = (0..hba_mem.ports.len()) .filter(|&i| pi & 1 << i as i32 == 1 << i as i32) .filter_map(|i| { let port = unsafe { &mut *hba_mem.ports.as_mut_ptr().add(i) }; let port_type = port.probe(); info!("{}-{}: {:?}", name, i, port_type); - let disk: Option> = match port_type { + let disk: Option = match port_type { HbaPortType::SATA => match DiskATA::new(i, port) { - Ok(disk) => Some(Box::new(disk)), + Ok(disk) => Some(AnyDisk::Ata(disk)), Err(err) => { error!("{}: {}", i, err); None } }, HbaPortType::SATAPI => match DiskATAPI::new(i, port) { - Ok(disk) => Some(Box::new(disk)), + Ok(disk) => Some(AnyDisk::Atapi(disk)), Err(err) => { error!("{}: {}", i, err); None diff --git a/storage/ahcid/src/main.rs b/storage/ahcid/src/main.rs index 9ea8964dd1..d6b9f83d1c 100644 --- a/storage/ahcid/src/main.rs +++ b/storage/ahcid/src/main.rs @@ -1,14 +1,11 @@ #![cfg_attr(target_arch = "aarch64", feature(stdsimd))] // Required for yield instruction -extern crate byteorder; -extern crate syscall; - use std::io::{Read, Write}; use std::os::fd::AsRawFd; use std::usize; use common::io::Io; -use driver_block::DiskScheme; +use driver_block::{DiskScheme, ExecutorTrait, FuturesExecutor}; use event::{EventFlags, RawEventQueue}; use pcid_interface::PciFunctionHandle; @@ -54,6 +51,7 @@ fn daemon(daemon: redox_daemon::Daemon) -> ! { .enumerate() .map(|(i, disk)| (i as u32, disk)) .collect(), + &FuturesExecutor, ); let mut irq_file = irq.irq_handle("ahcid"); @@ -75,7 +73,7 @@ fn daemon(daemon: redox_daemon::Daemon) -> ! { for event in event_queue { let event = event.unwrap(); if event.fd == scheme.event_handle().raw() { - scheme.tick().unwrap(); + FuturesExecutor.block_on(scheme.tick()).unwrap(); } else if event.fd == irq_fd { let mut irq = [0; 8]; if irq_file @@ -100,7 +98,7 @@ fn daemon(daemon: redox_daemon::Daemon) -> ! { .write(&irq) .expect("ahcid: failed to write irq file"); - scheme.tick().unwrap(); + FuturesExecutor.block_on(scheme.tick()).unwrap(); } } } else { diff --git a/storage/driver-block/Cargo.toml b/storage/driver-block/Cargo.toml index 00e10150bb..da969a41c8 100644 --- a/storage/driver-block/Cargo.toml +++ b/storage/driver-block/Cargo.toml @@ -4,8 +4,14 @@ version = "0.1.0" edition = "2021" [dependencies] +executor = { path = "../../executor" } partitionlib = { path = "../partitionlib" } libredox = "0.1.3" -redox_syscall = "0.5" -redox-scheme = "0.4" +log = "0.4" + +# TODO: migrate virtio to our executor +futures = { version = "0.3.28", features = ["executor"] } + +redox_syscall = { version = "0.5", features = ["std"] } +redox-scheme = "0.5" diff --git a/storage/driver-block/src/lib.rs b/storage/driver-block/src/lib.rs index 6efe0f5c2c..0363228587 100644 --- a/storage/driver-block/src/lib.rs +++ b/storage/driver-block/src/lib.rs @@ -1,20 +1,23 @@ use std::cmp; +use std::future::{Future, IntoFuture}; use std::io::{self, Read, Seek, SeekFrom}; use std::collections::BTreeMap; use std::convert::TryFrom; use std::fmt::Write; use std::str; +use std::task::Poll; +use executor::LocalExecutor; use libredox::Fd; use partitionlib::{LogicalBlockSize, PartitionTable}; -use redox_scheme::{ - CallRequest, CallerCtx, OpenResult, RequestKind, SchemeBlock, SignalBehavior, Socket, -}; +use redox_scheme::scheme::SchemeAsync; +use redox_scheme::{CallerCtx, OpenResult, RequestKind, Response, SignalBehavior, Socket}; +use syscall::dirent::DirentBuf; use syscall::schemev2::NewFdFlags; use syscall::{ - Error, Result, Stat, EACCES, EAGAIN, EBADF, EINVAL, EISDIR, ENOENT, ENOLCK, EOVERFLOW, - MODE_DIR, MODE_FILE, O_DIRECTORY, O_STAT, + Error, Result, Stat, EACCES, EAGAIN, EBADF, EINTR, EINVAL, EISDIR, ENOENT, ENOLCK, EOPNOTSUPP, + EOVERFLOW, EWOULDBLOCK, MODE_DIR, MODE_FILE, O_DIRECTORY, O_STAT, }; /// Split the read operation into a series of block reads. @@ -71,8 +74,8 @@ pub trait Disk { // These operate on a whole multiple of the block size // FIXME maybe only operate on a single block worth of data? - fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result>; - fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result>; + async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result; + async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result; } impl Disk for Box { @@ -84,12 +87,12 @@ impl Disk for Box { (**self).size() } - fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result> { - (**self).read(block, buffer) + async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result { + (**self).read(block, buffer).await } - fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result> { - (**self).write(block, buffer) + async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result { + (**self).write(block, buffer).await } } @@ -99,18 +102,19 @@ pub struct DiskWrapper { } impl DiskWrapper { - pub fn pt(disk: &mut T) -> Option { + pub fn pt(disk: &mut T, executor: &impl ExecutorTrait) -> Option { let bs = match disk.block_size() { 512 => LogicalBlockSize::Lb512, 4096 => LogicalBlockSize::Lb4096, _ => return None, }; - struct Device<'a> { - disk: &'a mut dyn Disk, + struct Device<'a, D: Disk, E: ExecutorTrait> { + disk: &'a mut D, + executor: &'a E, offset: u64, } - impl<'a> Seek for Device<'a> { + impl<'a, D: Disk, E: ExecutorTrait> Seek for Device<'a, D, E> { fn seek(&mut self, from: SeekFrom) -> io::Result { let size = i64::try_from(self.disk.size()).or(Err(io::Error::new( io::ErrorKind::Other, @@ -129,7 +133,7 @@ impl DiskWrapper { } } // TODO: Perhaps this impl should be used in the rest of the scheme. - impl<'a> Read for Device<'a> { + impl<'a, D: Disk, E: ExecutorTrait> Read for Device<'a, D, E> { fn read(&mut self, buf: &mut [u8]) -> io::Result { let blksize = self.disk.block_size(); let size_in_blocks = self.disk.size() / u64::from(blksize); @@ -141,17 +145,9 @@ impl DiskWrapper { return Err(io::Error::from_raw_os_error(syscall::EOVERFLOW)); } loop { - match disk.read(block, block_bytes) { - Ok(Some(bytes)) => { - assert_eq!(bytes, block_bytes.len()); - return Ok(()); - } - Ok(None) => { - std::thread::yield_now(); - continue; - } - Err(err) => return Err(io::Error::from_raw_os_error(err.errno)), - } + let bytes = self.executor.block_on(disk.read(block, block_bytes))?; + assert_eq!(bytes, block_bytes.len()); + return Ok(()); } }; let bytes_read = block_read(self.offset, blksize, buf, read_block)?; @@ -161,14 +157,21 @@ impl DiskWrapper { } } - partitionlib::get_partitions(&mut Device { disk, offset: 0 }, bs) - .ok() - .flatten() + partitionlib::get_partitions( + &mut Device { + disk, + offset: 0, + executor, + }, + bs, + ) + .ok() + .flatten() } - pub fn new(mut disk: T) -> Self { + pub fn new(mut disk: T, executor: &impl ExecutorTrait) -> Self { Self { - pt: Self::pt(&mut disk), + pt: Self::pt(&mut disk, executor), disk, } } @@ -189,12 +192,12 @@ impl DiskWrapper { self.disk.size() } - pub fn read( + pub async fn read( &mut self, part_num: Option, block: u64, buf: &mut [u8], - ) -> syscall::Result> { + ) -> syscall::Result { if buf.len() as u64 % u64::from(self.disk.block_size()) != 0 { return Err(Error::new(EINVAL)); } @@ -214,18 +217,18 @@ impl DiskWrapper { let abs_block = part.start_lba + block; - self.disk.read(abs_block, buf) + self.disk.read(abs_block, buf).await } else { - self.disk.read(block, buf) + self.disk.read(block, buf).await } } - pub fn write( + pub async fn write( &mut self, part_num: Option, block: u64, buf: &[u8], - ) -> syscall::Result> { + ) -> syscall::Result { if buf.len() as u64 % u64::from(self.disk.block_size()) != 0 { return Err(Error::new(EINVAL)); } @@ -245,9 +248,9 @@ impl DiskWrapper { let abs_block = part.start_lba + block; - self.disk.write(abs_block, buf) + self.disk.write(abs_block, buf).await } else { - self.disk.write(block, buf) + self.disk.write(block, buf).await } } } @@ -264,11 +267,48 @@ pub struct DiskScheme { disks: BTreeMap>, handles: BTreeMap, next_id: usize, - blocked: Vec, +} + +pub trait ExecutorTrait { + fn block_on<'a, O: 'a>(&self, fut: impl IntoFuture + 'a) -> O; +} +impl ExecutorTrait for LocalExecutor { + fn block_on<'a, O: 'a>(&self, fut: impl IntoFuture + 'a) -> O { + LocalExecutor::block_on(self, fut) + } +} +#[deprecated = "use custom executor"] +pub struct FuturesExecutor; + +#[allow(deprecated)] +impl ExecutorTrait for FuturesExecutor { + fn block_on<'a, O: 'a>(&self, fut: impl IntoFuture + 'a) -> O { + futures::executor::block_on(fut.into_future()) + } +} +pub struct TrivialExecutor; +impl ExecutorTrait for TrivialExecutor { + fn block_on<'a, O: 'a>(&self, fut: impl IntoFuture + 'a) -> O { + let mut fut = std::pin::pin!(fut.into_future()); + let mut cx = std::task::Context::from_waker(std::task::Waker::noop()); + loop { + match fut.as_mut().poll(&mut cx) { + Poll::Ready(v) => return v, + Poll::Pending => { + log::warn!("TrivialExecutor: future wasn't trivial"); + continue; + } + } + } + } } impl DiskScheme { - pub fn new(scheme_name: String, disks: BTreeMap) -> Self { + pub fn new( + scheme_name: String, + disks: BTreeMap, + executor: &impl ExecutorTrait, + ) -> Self { assert!(scheme_name.starts_with("disk")); let socket = Socket::nonblock(&scheme_name).expect("failed to create disk scheme"); @@ -277,11 +317,10 @@ impl DiskScheme { socket, disks: disks .into_iter() - .map(|(k, disk)| (k, DiskWrapper::new(disk))) + .map(|(k, disk)| (k, DiskWrapper::new(disk, executor))) .collect(), next_id: 0, handles: BTreeMap::new(), - blocked: vec![], } } @@ -291,53 +330,45 @@ impl DiskScheme { /// Process pending and new requests. /// - /// This needs to be called each time there is a new event on the scheme - /// file and each time a read or write operation has completed. - // FIXME maybe split into one method for events on the scheme fd and one - // to call when an irq is received to indicate that blocked packets can - // be processed. - pub fn tick(&mut self) -> io::Result<()> { - // Handle any blocked requests - let mut i = 0; - while i < self.blocked.len() { - if let Some(resp) = self.blocked[i].handle_scheme_block(self) { - self.socket - .write_response(resp, SignalBehavior::Restart) - .expect("driver-block: failed to write scheme"); - self.blocked.remove(i); - } else { - i += 1; - } - } - + /// This needs to be called each time there is a new event on the scheme. + pub async fn tick(&mut self) -> io::Result<()> { // Handle new scheme requests loop { - let request = match self.socket.next_request(SignalBehavior::Restart) { + let request = match self.socket.next_request(SignalBehavior::Interrupt) { Ok(Some(request)) => request, Ok(None) => { // Scheme likely got unmounted + // TODO: return this to caller instead std::process::exit(0); } - Err(err) if err.errno == EAGAIN => break, + Err(error) if error.errno == EWOULDBLOCK || error.errno == EAGAIN => break, + Err(err) if err.errno == EINTR => continue, Err(err) => return Err(err.into()), }; - match request.kind() { + let response = match request.kind() { RequestKind::Call(call_request) => { - if let Some(resp) = call_request.handle_scheme_block(self) { - self.socket.write_response(resp, SignalBehavior::Restart)?; - } else { - self.blocked.push(call_request); - } + // TODO: Spawn a separate task for each scheme call. This would however require the + // use of a smarter buffer pool (or direct IO, or a buffer per fd) in order to do + // parallel IO. It might also require async-aware locks so that a close() is + // correctly ordered wrt IO on the same fd. + call_request.handle_async(self).await + } + RequestKind::SendFd(sendfd_request) => Response::err(EOPNOTSUPP, sendfd_request), + RequestKind::Cancellation(_cancellation_request) => { + // FIXME implement cancellation + continue; + } + RequestKind::MsyncMsg | RequestKind::MunmapMsg | RequestKind::MmapMsg => { + unreachable!() } RequestKind::OnClose { id } => { self.on_close(id); + continue; } - RequestKind::Cancellation(_cancellation_request) => { - // FIXME implement cancellation - } - _ => {} - } + }; + self.socket + .write_response(response, SignalBehavior::Restart)?; } Ok(()) @@ -373,13 +404,8 @@ impl DiskScheme { } } -impl SchemeBlock for DiskScheme { - fn xopen( - &mut self, - path_str: &str, - flags: usize, - ctx: &CallerCtx, - ) -> Result> { +impl SchemeAsync for DiskScheme { + async fn open(&mut self, path_str: &str, flags: usize, ctx: &CallerCtx) -> Result { if ctx.uid != 0 { return Err(Error::new(EACCES)); } @@ -446,18 +472,27 @@ impl SchemeBlock for DiskScheme { let id = self.next_id; self.next_id += 1; self.handles.insert(id, handle); - Ok(Some(OpenResult::ThisScheme { + Ok(OpenResult::ThisScheme { number: id, flags: NewFdFlags::POSITIONED, - })) + }) + } + async fn getdents<'buf>( + &mut self, + _id: usize, + _buf: DirentBuf<&'buf mut [u8]>, + _opaque_offset: u64, + ) -> Result> { + // TODO + Err(Error::new(EOPNOTSUPP)) } - fn fstat(&mut self, id: usize, stat: &mut Stat) -> Result> { + async fn fstat(&mut self, id: usize, stat: &mut Stat, _ctx: &CallerCtx) -> Result<()> { match *self.handles.get(&id).ok_or(Error::new(EBADF))? { Handle::List(ref data) => { stat.st_mode = MODE_DIR; stat.st_size = data.len() as u64; - Ok(Some(0)) + Ok(()) } Handle::Disk(number) => { let disk = self.disks.get_mut(&number).ok_or(Error::new(EBADF))?; @@ -465,7 +500,7 @@ impl SchemeBlock for DiskScheme { stat.st_blocks = disk.disk().size() / u64::from(disk.block_size()); stat.st_blksize = disk.block_size(); stat.st_size = disk.size(); - Ok(Some(0)) + Ok(()) } Handle::Partition(disk_num, part_num) => { let disk = self.disks.get_mut(&disk_num).ok_or(Error::new(EBADF))?; @@ -480,18 +515,19 @@ impl SchemeBlock for DiskScheme { stat.st_size = part.size * u64::from(disk.block_size()); stat.st_blocks = part.size; stat.st_blksize = disk.block_size(); - Ok(Some(0)) + Ok(()) } } } - fn fpath(&mut self, id: usize, buf: &mut [u8]) -> Result> { + async fn fpath(&mut self, id: usize, buf: &mut [u8], _ctx: &CallerCtx) -> Result { let handle = self.handles.get(&id).ok_or(Error::new(EBADF))?; let mut i = 0; let scheme_name = self.scheme_name.as_bytes(); let mut j = 0; + // TODO: copy_from_slice while i < buf.len() && j < scheme_name.len() { buf[i] = scheme_name[j]; i += 1; @@ -527,16 +563,17 @@ impl SchemeBlock for DiskScheme { } } - Ok(Some(i)) + Ok(i) } - fn read( + async fn read( &mut self, id: usize, buf: &mut [u8], offset: u64, _fcntl_flags: u32, - ) -> Result> { + _ctx: &CallerCtx, + ) -> Result { match *self.handles.get_mut(&id).ok_or(Error::new(EBADF))? { Handle::List(ref handle) => { let src = usize::try_from(offset) @@ -545,70 +582,69 @@ impl SchemeBlock for DiskScheme { .unwrap_or(&[]); let count = core::cmp::min(src.len(), buf.len()); buf[..count].copy_from_slice(&src[..count]); - Ok(Some(count)) + Ok(count) } Handle::Disk(number) => { let disk = self.disks.get_mut(&number).ok_or(Error::new(EBADF))?; let block = offset / u64::from(disk.block_size()); - disk.read(None, block, buf) + disk.read(None, block, buf).await } Handle::Partition(disk_num, part_num) => { let disk = self.disks.get_mut(&disk_num).ok_or(Error::new(EBADF))?; let block = offset / u64::from(disk.block_size()); - disk.read(Some(part_num as usize), block, buf) + disk.read(Some(part_num as usize), block, buf).await } } } - fn write( + async fn write( &mut self, id: usize, buf: &[u8], offset: u64, _fcntl_flags: u32, - ) -> Result> { + _ctx: &CallerCtx, + ) -> Result { match *self.handles.get_mut(&id).ok_or(Error::new(EBADF))? { Handle::List(_) => Err(Error::new(EBADF)), Handle::Disk(number) => { let disk = self.disks.get_mut(&number).ok_or(Error::new(EBADF))?; let block = offset / u64::from(disk.block_size()); - disk.write(None, block, buf) + disk.write(None, block, buf).await } Handle::Partition(disk_num, part_num) => { let disk = self.disks.get_mut(&disk_num).ok_or(Error::new(EBADF))?; let block = offset / u64::from(disk.block_size()); - disk.write(Some(part_num as usize), block, buf) + disk.write(Some(part_num as usize), block, buf).await } } } - fn fsize(&mut self, id: usize) -> Result> { - Ok(Some( - match *self.handles.get_mut(&id).ok_or(Error::new(EBADF))? { - Handle::List(ref handle) => handle.len() as u64, - Handle::Disk(number) => { - let disk = self.disks.get_mut(&number).ok_or(Error::new(EBADF))?; - disk.size() - } - Handle::Partition(disk_num, part_num) => { - let disk = self.disks.get_mut(&disk_num).ok_or(Error::new(EBADF))?; - let part = disk - .pt - .as_ref() - .ok_or(Error::new(EBADF))? - .partitions - .get(part_num as usize) - .ok_or(Error::new(EBADF))?; + async fn fsize(&mut self, id: usize, _ctx: &CallerCtx) -> Result { + Ok(match *self.handles.get_mut(&id).ok_or(Error::new(EBADF))? { + Handle::List(ref handle) => handle.len() as u64, + Handle::Disk(number) => { + let disk = self.disks.get_mut(&number).ok_or(Error::new(EBADF))?; + disk.size() + } + Handle::Partition(disk_num, part_num) => { + let disk = self.disks.get_mut(&disk_num).ok_or(Error::new(EBADF))?; + let part = disk + .pt + .as_ref() + .ok_or(Error::new(EBADF))? + .partitions + .get(part_num as usize) + .ok_or(Error::new(EBADF))?; - part.size * u64::from(disk.block_size()) - } - }, - )) + part.size * u64::from(disk.block_size()) + } + }) } } -impl DiskScheme { - fn on_close(&mut self, id: usize) { - self.handles.remove(&id); +impl DiskScheme { + pub fn on_close(&mut self, id: usize) { + let _ = self.handles.remove(&id); } } diff --git a/storage/ided/src/ide.rs b/storage/ided/src/ide.rs index 9748b46b93..0a14bcc99a 100644 --- a/storage/ided/src/ide.rs +++ b/storage/ided/src/ide.rs @@ -177,7 +177,8 @@ impl Disk for AtaDisk { self.size } - fn read(&mut self, start_block: u64, buffer: &mut [u8]) -> Result> { + // NOTE: not async + async fn read(&mut self, start_block: u64, buffer: &mut [u8]) -> Result { let mut count = 0; for chunk in buffer.chunks_mut(65536) { let block = start_block + (count as u64) / 512; @@ -314,10 +315,11 @@ impl Disk for AtaDisk { count += chunk.len(); } - Ok(Some(count)) + Ok(count) } - fn write(&mut self, start_block: u64, buffer: &[u8]) -> Result> { + // NOTE: not async + async fn write(&mut self, start_block: u64, buffer: &[u8]) -> Result { let mut count = 0; for chunk in buffer.chunks(65536) { let block = start_block + (count as u64) / 512; @@ -462,6 +464,6 @@ impl Disk for AtaDisk { count += chunk.len(); } - Ok(Some(count)) + Ok(count) } } diff --git a/storage/ided/src/main.rs b/storage/ided/src/main.rs index c7dd8f8746..d3509a1ba8 100644 --- a/storage/ided/src/main.rs +++ b/storage/ided/src/main.rs @@ -1,5 +1,5 @@ use common::io::Io as _; -use driver_block::{Disk, DiskScheme}; +use driver_block::{Disk, DiskScheme, ExecutorTrait, FuturesExecutor}; use event::{EventFlags, RawEventQueue}; use libredox::flag; use log::{error, info}; @@ -61,7 +61,28 @@ fn daemon(daemon: redox_daemon::Daemon) -> ! { Arc::new(Mutex::new(primary)), Arc::new(Mutex::new(secondary)), ]; - let mut disks: Vec> = Vec::new(); + enum AnyDisk { + Ata(AtaDisk), + } + impl Disk for AnyDisk { + fn block_size(&self) -> u32 { + let AnyDisk::Ata(a) = self; + a.block_size() + } + fn size(&self) -> u64 { + let AnyDisk::Ata(a) = self; + a.size() + } + async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result { + let AnyDisk::Ata(a) = self; + a.write(block, buffer).await + } + async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result { + let AnyDisk::Ata(a) = self; + a.read(block, buffer).await + } + } + let mut disks: Vec = Vec::new(); for (chan_i, chan_lock) in chans.iter().enumerate() { let mut chan = chan_lock.lock().unwrap(); @@ -174,7 +195,7 @@ fn daemon(daemon: redox_daemon::Daemon) -> ! { println!(" DMA: {}", dma); println!(" {}-bit LBA", lba_bits); - disks.push(Box::new(AtaDisk { + disks.push(AnyDisk::Ata(AtaDisk { chan: chan_lock.clone(), chan_i, dev, @@ -194,6 +215,9 @@ fn daemon(daemon: redox_daemon::Daemon) -> ! { .enumerate() .map(|(i, disk)| (i as u32, disk)) .collect(), + // TODO: Should ided just use TrivialExecutor or would it be valuable to actually use a + // real executor? + &FuturesExecutor, ); let primary_irq_fd = libredox::call::open( @@ -233,7 +257,7 @@ fn daemon(daemon: redox_daemon::Daemon) -> ! { for event in event_queue { let event = event.unwrap(); if event.fd == scheme.event_handle().raw() { - scheme.tick().unwrap(); + FuturesExecutor.block_on(scheme.tick()).unwrap(); } else if event.fd == primary_irq_fd { let mut irq = [0; 8]; if primary_irq_file @@ -248,7 +272,7 @@ fn daemon(daemon: redox_daemon::Daemon) -> ! { .write(&irq) .expect("ided: failed to write irq file"); - scheme.tick().unwrap(); + FuturesExecutor.block_on(scheme.tick()).unwrap(); } } else if event.fd == secondary_irq_fd { let mut irq = [0; 8]; @@ -264,7 +288,7 @@ fn daemon(daemon: redox_daemon::Daemon) -> ! { .write(&irq) .expect("ided: failed to write irq file"); - scheme.tick().unwrap(); + FuturesExecutor.block_on(scheme.tick()).unwrap(); } } else { error!("Unknown event {}", event.fd); diff --git a/storage/lived/src/main.rs b/storage/lived/src/main.rs index f4c2ac15ff..b53452e779 100644 --- a/storage/lived/src/main.rs +++ b/storage/lived/src/main.rs @@ -8,6 +8,7 @@ use std::fs::File; use std::os::fd::AsRawFd; use driver_block::{Disk, DiskScheme}; +use driver_block::{ExecutorTrait, TrivialExecutor}; use libredox::call::MmapArgs; use libredox::flag; @@ -86,7 +87,7 @@ impl Disk for LiveDisk { self.the_data.len() as u64 } - fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result> { + async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result { let block = block as usize; let block_size = self.block_size() as usize; if block * block_size + buffer.len() > self.size() as usize { @@ -94,10 +95,10 @@ impl Disk for LiveDisk { } buffer .copy_from_slice(&self.the_data[block * block_size..block * block_size + buffer.len()]); - Ok(Some(block_size)) + Ok(block_size) } - fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result> { + async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result { let block = block as usize; let block_size = self.block_size() as usize; if block * block_size + buffer.len() > self.size() as usize { @@ -105,7 +106,7 @@ impl Disk for LiveDisk { } self.the_data[block * block_size..block * block_size + buffer.len()] .copy_from_slice(buffer); - Ok(Some(block_size)) + Ok(block_size) } } @@ -128,6 +129,7 @@ fn main() -> anyhow::Result<()> { std::process::exit(1) }), )]), + &TrivialExecutor, ); libredox::call::setrens(0, 0).expect("nvmed: failed to enter null namespace"); @@ -144,7 +146,7 @@ fn main() -> anyhow::Result<()> { for event in event_queue { match event.unwrap().user_data { - Event::Scheme => scheme.tick().unwrap(), + Event::Scheme => TrivialExecutor.block_on(scheme.tick()).unwrap(), } } diff --git a/storage/nvmed/Cargo.toml b/storage/nvmed/Cargo.toml index 3650871a87..948f30a970 100644 --- a/storage/nvmed/Cargo.toml +++ b/storage/nvmed/Cargo.toml @@ -4,21 +4,21 @@ version = "0.1.0" edition = "2021" [dependencies] -arrayvec = "0.5" -bitflags = "1" -crossbeam-channel = "0.4" -futures = "0.3" +arrayvec = "0.7" +bitflags = "2" +libredox = "0.1.3" log = "0.4" +parking_lot = "0.12.1" redox-daemon = "0.1" +redox_event = "0.4.1" redox_syscall = { version = "0.5", features = ["std"] } -redox_event = "0.4" smallvec = "1" +executor = { path = "../../executor" } common = { path = "../../common" } driver-block = { path = "../driver-block" } +partitionlib = { path = "../partitionlib" } pcid = { path = "../../pcid" } -libredox = "0.1.3" [features] -default = ["async"] -async = [] +default = [] diff --git a/storage/nvmed/src/main.rs b/storage/nvmed/src/main.rs index d53189d1f2..efe390f967 100644 --- a/storage/nvmed/src/main.rs +++ b/storage/nvmed/src/main.rs @@ -1,7 +1,9 @@ #![cfg_attr(target_arch = "aarch64", feature(stdarch_arm_hints))] // Required for yield instruction #![cfg_attr(target_arch = "riscv64", feature(riscv_ext_intrinsics))] // Required for pause instruction +use std::cell::RefCell; use std::ptr::NonNull; +use std::rc::Rc; use std::sync::Arc; use std::{slice, usize}; @@ -140,23 +142,26 @@ fn get_int_method( } } -struct NvmeDisk(Arc, NvmeNamespace); +struct NvmeDisk { + nvme: Arc, + ns: NvmeNamespace, +} impl Disk for NvmeDisk { fn block_size(&self) -> u32 { - self.1.block_size.try_into().unwrap() + self.ns.block_size.try_into().unwrap() } fn size(&self) -> u64 { - self.1.blocks * self.1.block_size + self.ns.blocks * self.ns.block_size } - fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result> { - self.0.namespace_read(self.1, block, buffer) + async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result { + self.nvme.namespace_read(&self.ns, block, buffer).await } - fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result> { - self.0.namespace_write(self.1, block, buffer) + async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result { + self.nvme.namespace_write(&self.ns, block, buffer).await } } @@ -181,65 +186,67 @@ fn daemon(daemon: redox_daemon::Daemon) -> ! { let address = unsafe { pcid_handle.map_bar(0).ptr }; - daemon.ready().expect("nvmed: failed to signal readiness"); - - let (reactor_sender, reactor_receiver) = crossbeam_channel::unbounded(); let (interrupt_method, interrupt_sources) = get_int_method(&mut pcid_handle, &pci_config.func) .expect("nvmed: failed to find a suitable interrupt method"); - let mut nvme = Nvme::new( - address.as_ptr() as usize, - interrupt_method, - pcid_handle, - reactor_sender, - ) - .expect("nvmed: failed to allocate driver data"); + let mut nvme = Nvme::new(address.as_ptr() as usize, interrupt_method, pcid_handle) + .expect("nvmed: failed to allocate driver data"); + unsafe { nvme.init() } log::debug!("Finished base initialization"); let nvme = Arc::new(nvme); - #[cfg(feature = "async")] - let reactor_thread = nvme::cq_reactor::start_cq_reactor_thread( - Arc::clone(&nvme), - interrupt_sources, - reactor_receiver, - ); - let namespaces = nvme.init_with_queues(); - let event_queue = event::EventQueue::new().unwrap(); - event::user_data! { - enum Event { - Scheme, - } + let executor = { + let (intx, (iv, irq_handle)) = match interrupt_sources { + InterruptSources::Msi(mut vectors) => ( + false, + vectors.pop_first().map(|(a, b)| (u16::from(a), b)).unwrap(), + ), + InterruptSources::MsiX(mut vectors) => (false, vectors.pop_first().unwrap()), + InterruptSources::Intx(file) => (true, (0, file)), + }; + nvme::executor::init(Arc::clone(&nvme), iv, intx, irq_handle) }; + let namespaces = executor.block_on(nvme.init_with_queues()); + log::debug!("Initialized!"); - let mut scheme = DiskScheme::new( + let scheme = Rc::new(RefCell::new(DiskScheme::new( scheme_name, namespaces .into_iter() - .map(|(k, ns)| (k, NvmeDisk(nvme.clone(), ns))) + .map(|(k, ns)| { + ( + k, + NvmeDisk { + nvme: nvme.clone(), + ns, + }, + ) + }) .collect(), - ); + &*executor, + ))); + daemon.ready().expect("nvmed: failed to signal readiness"); + + let mut scheme_events = Box::pin(executor.register_external_event( + scheme.borrow().event_handle().raw(), + event::EventFlags::READ, + )); libredox::call::setrens(0, 0).expect("nvmed: failed to enter null namespace"); - event_queue - .subscribe( - scheme.event_handle().raw(), - Event::Scheme, - event::EventFlags::READ, - ) - .unwrap(); + log::info!("Starting to listen for scheme events"); - for event in event_queue { - match event.unwrap().user_data { - Event::Scheme => scheme.tick().unwrap(), + executor.block_on(async { + loop { + log::trace!("new event iteration"); + if let Err(err) = scheme.borrow_mut().tick().await { + log::error!("scheme error: {err}"); + } + let _ = scheme_events.as_mut().next().await; } - } + }); //TODO: destroy NVMe stuff - #[cfg(feature = "async")] - reactor_thread - .join() - .expect("nvmed: failed to join reactor thread"); std::process::exit(0); } diff --git a/storage/nvmed/src/nvme/executor.rs b/storage/nvmed/src/nvme/executor.rs new file mode 100644 index 0000000000..6242fa98cb --- /dev/null +++ b/storage/nvmed/src/nvme/executor.rs @@ -0,0 +1,82 @@ +use std::cell::RefCell; +use std::fs::File; +use std::rc::Rc; +use std::sync::Arc; + +use executor::{Hardware, LocalExecutor}; + +use super::{CmdId, CqId, Nvme, NvmeCmd, NvmeComp, SqId}; + +pub struct NvmeHw; + +impl Hardware for NvmeHw { + type Iv = u16; + type Sqe = NvmeCmd; + type Cqe = NvmeComp; + type CmdId = CmdId; + type CqId = CqId; + type SqId = SqId; + type GlobalCtxt = Arc; + + fn mask_vector(ctxt: &Arc, iv: Self::Iv) { + ctxt.set_vector_masked(iv, true) + } + fn unmask_vector(ctxt: &Arc, iv: Self::Iv) { + ctxt.set_vector_masked(iv, false) + } + fn set_sqe_cmdid(sqe: &mut NvmeCmd, id: CmdId) { + sqe.cid = id; + } + fn get_cqe_cmdid(cqe: &Self::Cqe) -> Self::CmdId { + cqe.cid + } + fn vtable() -> &'static std::task::RawWakerVTable { + &VTABLE + } + fn current() -> std::rc::Rc> { + THE_EXECUTOR.with(|exec| Rc::clone(exec.borrow().as_ref().unwrap())) + } + fn try_submit( + nvme: &Arc, + sq_id: Self::SqId, + success: impl FnOnce(Self::CmdId) -> Self::Sqe, + fail: impl FnOnce(), + ) -> Option<(Self::CqId, Self::CmdId)> { + let ctxt = nvme.cur_thread_ctxt(); + let ctxt = ctxt.lock(); + + nvme.try_submit_raw(&*ctxt, sq_id, success, fail) + } + fn poll_cqes(nvme: &Arc, mut handle: impl FnMut(Self::CqId, Self::Cqe)) { + let ctxt = nvme.cur_thread_ctxt(); + let ctxt = ctxt.lock(); + + for (sq_cq_id, (sq, cq)) in ctxt.queues.borrow_mut().iter_mut() { + while let Some((new_head, cqe)) = cq.complete() { + unsafe { + nvme.completion_queue_head(*sq_cq_id, new_head); + } + sq.head = cqe.sq_head; + log::trace!("new head {new_head} cqe {cqe:?}"); + handle(*sq_cq_id, cqe); + } + } + } + fn sq_cq(_ctxt: &Arc, id: Self::CqId) -> Self::SqId { + id + } +} + +static VTABLE: std::task::RawWakerVTable = executor::vtable::(); + +thread_local! { + static THE_EXECUTOR: RefCell>>> = RefCell::new(None); +} + +pub type NvmeExecutor = LocalExecutor; + +pub fn init(nvme: Arc, iv: u16, intx: bool, irq_handle: File) -> Rc> { + let this = Rc::new(executor::init_raw(nvme, iv, intx, irq_handle)); + THE_EXECUTOR.with(|exec| *exec.borrow_mut() = Some(Rc::clone(&this))); + this +} diff --git a/storage/nvmed/src/nvme/identify.rs b/storage/nvmed/src/nvme/identify.rs index 93f221039d..05e5b9b2b6 100644 --- a/storage/nvmed/src/nvme/identify.rs +++ b/storage/nvmed/src/nvme/identify.rs @@ -151,14 +151,16 @@ impl LbaFormat { impl Nvme { /// Returns the serial number, model, and firmware, in that order. - pub fn identify_controller(&self) { + pub async fn identify_controller(&self) { // TODO: Use same buffer let data: Dma = unsafe { Dma::zeroed().unwrap().assume_init() }; // println!(" - Attempting to identify controller"); - let comp = self.submit_and_complete_admin_command(|cid| { - NvmeCmd::identify_controller(cid, data.physical()) - }); + let comp = self + .submit_and_complete_admin_command(|cid| { + NvmeCmd::identify_controller(cid, data.physical()) + }) + .await; log::trace!("Completion: {:?}", comp); // println!(" - Dumping identify controller"); @@ -178,30 +180,34 @@ impl Nvme { firmware, ); } - pub fn identify_namespace_list(&self, base: u32) -> Vec { + pub async fn identify_namespace_list(&self, base: u32) -> Vec { // TODO: Use buffer let data: Dma<[u32; 1024]> = unsafe { Dma::zeroed().unwrap().assume_init() }; // println!(" - Attempting to retrieve namespace ID list"); - let comp = self.submit_and_complete_admin_command(|cid| { - NvmeCmd::identify_namespace_list(cid, data.physical(), base) - }); + let comp = self + .submit_and_complete_admin_command(|cid| { + NvmeCmd::identify_namespace_list(cid, data.physical(), base) + }) + .await; log::trace!("Completion2: {:?}", comp); // println!(" - Dumping namespace ID list"); data.iter().copied().take_while(|&nsid| nsid != 0).collect() } - pub fn identify_namespace(&self, nsid: u32) -> NvmeNamespace { + pub async fn identify_namespace(&self, nsid: u32) -> NvmeNamespace { //TODO: Use buffer let data: Dma = unsafe { Dma::zeroed().unwrap().assume_init() }; - // println!(" - Attempting to identify namespace {}", nsid); - let comp = self.submit_and_complete_admin_command(|cid| { - NvmeCmd::identify_namespace(cid, data.physical(), nsid) - }); + log::debug!("Attempting to identify namespace {nsid}"); + let comp = self + .submit_and_complete_admin_command(|cid| { + NvmeCmd::identify_namespace(cid, data.physical(), nsid) + }) + .await; - // println!(" - Dumping identify namespace"); + log::debug!("Dumping identify namespace"); let size = data.size_in_blocks(); let capacity = data.capacity_in_blocks(); diff --git a/storage/nvmed/src/nvme/mod.rs b/storage/nvmed/src/nvme/mod.rs index 93a1c03b9b..513814c9d0 100644 --- a/storage/nvmed/src/nvme/mod.rs +++ b/storage/nvmed/src/nvme/mod.rs @@ -1,11 +1,12 @@ -use std::collections::BTreeMap; +use std::cell::RefCell; +use std::collections::{BTreeMap, HashMap}; use std::convert::TryFrom; use std::fs::File; -use std::sync::atomic::{AtomicU16, AtomicU64}; -use std::sync::{Mutex, RwLock}; +use std::iter; +use std::sync::atomic::AtomicU16; +use std::sync::Arc; -use crossbeam_channel::Sender; -use smallvec::{smallvec, SmallVec}; +use parking_lot::{Mutex, ReentrantMutex, RwLock}; use common::io::{Io, Mmio}; use syscall::error::{Error, Result, EIO}; @@ -13,11 +14,11 @@ use syscall::error::{Error, Result, EIO}; use common::dma::Dma; pub mod cmd; -pub mod cq_reactor; +pub mod executor; pub mod identify; pub mod queues; -use self::cq_reactor::NotifReq; +use self::executor::NvmeExecutor; pub use self::queues::{NvmeCmd, NvmeCmdQueue, NvmeComp, NvmeCompQueue}; use pcid_interface::msi::{MsiInfo, MsixInfo, MsixTableEntry}; @@ -47,7 +48,7 @@ pub(crate) unsafe fn pause() { std::arch::riscv64::pause(); } -/// Used in conjunction with `InterruptMethod`, primarily by the CQ reactor. +/// Used in conjunction with `InterruptMethod`, primarily by the CQ executor. #[derive(Debug)] pub enum InterruptSources { MsiX(BTreeMap), @@ -183,28 +184,32 @@ pub type CmdId = u16; pub type AtomicCqId = AtomicU16; pub type AtomicSqId = AtomicU16; pub type AtomicCmdId = AtomicU16; +pub type Iv = u16; pub struct Nvme { interrupt_method: Mutex, pcid_interface: Mutex, regs: RwLock<&'static mut NvmeRegs>, - pub(crate) submission_queues: RwLock, CqId)>>, - pub(crate) completion_queues: - RwLock)>>>, + sq_ivs: RwLock>, + cq_ivs: RwLock>, // maps interrupt vectors with the completion queues they have - cqs_for_ivs: RwLock>>, - - buffer: Mutex>, // 2MB of buffer - buffer_prp: Mutex>, // 4KB of PRP for the buffer - reactor_sender: Sender, + thread_ctxts: RwLock>>>, next_sqid: AtomicSqId, next_cqid: AtomicCqId, - - next_avail_submission_epoch: AtomicU64, } + +pub struct ThreadCtxt { + buffer: RefCell>, // 2MB of buffer + buffer_prp: RefCell>, // 4KB of PRP for the buffer + + // Yes, technically NVME allows multiple submission queues to be mapped to the same completion + // queue, but we don't use that feature. + queues: RefCell>, +} + unsafe impl Send for Nvme {} unsafe impl Sync for Nvme {} @@ -213,8 +218,8 @@ pub enum FullSqHandling { /// Return an error immediately prior to posting the command. ErrorDirectly, - /// Tell the IRQ reactor that we want to be notified when a command on the same submission - /// queue has been completed. + /// Tell the executor that we want to be notified when a command on the same submission queue + /// has been completed. Wait, } @@ -223,29 +228,34 @@ impl Nvme { address: usize, interrupt_method: InterruptMethod, pcid_interface: PciFunctionHandle, - reactor_sender: Sender, ) -> Result { Ok(Nvme { regs: RwLock::new(unsafe { &mut *(address as *mut NvmeRegs) }), - submission_queues: RwLock::new( - std::iter::once((0u16, (Mutex::new(NvmeCmdQueue::new()?), 0u16))).collect(), + thread_ctxts: RwLock::new( + iter::once(( + 0_u16, + Arc::new(ReentrantMutex::new(ThreadCtxt { + buffer: RefCell::new(unsafe { Dma::zeroed()?.assume_init() }), + buffer_prp: RefCell::new(unsafe { Dma::zeroed()?.assume_init() }), + + queues: RefCell::new( + iter::once((0, (NvmeCmdQueue::new()?, NvmeCompQueue::new()?))) + .collect(), + ), + })), + )) + .collect(), ), - completion_queues: RwLock::new( - std::iter::once((0u16, Mutex::new((NvmeCompQueue::new()?, smallvec!(0))))) - .collect(), - ), - // map the zero interrupt vector (which according to the spec shall always point to the - // admin completion queue) to CQID 0 (admin completion queue) - cqs_for_ivs: RwLock::new(std::iter::once((0, smallvec!(0))).collect()), - buffer: Mutex::new(unsafe { Dma::zeroed()?.assume_init() }), - buffer_prp: Mutex::new(unsafe { Dma::zeroed()?.assume_init() }), + + cq_ivs: RwLock::new(iter::once((0, 0)).collect()), + sq_ivs: RwLock::new(iter::once((0, 0)).collect()), + interrupt_method: Mutex::new(interrupt_method), pcid_interface: Mutex::new(pcid_interface), - reactor_sender, - next_sqid: AtomicSqId::new(0), - next_cqid: AtomicCqId::new(0), - next_avail_submission_epoch: AtomicU64::new(0), + // TODO + next_sqid: AtomicSqId::new(2), + next_cqid: AtomicCqId::new(2), }) } /// Write to a doorbell register. @@ -255,13 +265,17 @@ impl Nvme { unsafe fn doorbell_write(&self, index: usize, value: u32) { use std::ops::DerefMut; - let mut regs_guard = self.regs.write().unwrap(); - let mut regs: &mut NvmeRegs = regs_guard.deref_mut(); + let mut regs_guard = self.regs.write(); + let regs: &mut NvmeRegs = regs_guard.deref_mut(); let dstrd = (regs.cap_high.read() & 0b1111) as usize; let addr = (regs as *mut NvmeRegs as usize) + 0x1000 + index * (4 << dstrd); (&mut *(addr as *mut Mmio)).write(value); } + fn cur_thread_ctxt(&self) -> Arc> { + // TODO: multi-threading + Arc::clone(self.thread_ctxts.read().get(&0).unwrap()) + } pub unsafe fn submission_queue_tail(&self, qid: u16, tail: u16) { self.doorbell_write(2 * (qid as usize), u32::from(tail)); @@ -272,15 +286,9 @@ impl Nvme { } pub unsafe fn init(&mut self) { - let mut buffer = self.buffer.get_mut().unwrap(); - let mut buffer_prp = self.buffer_prp.get_mut().unwrap(); - - for i in 0..buffer_prp.len() { - buffer_prp[i] = (buffer.physical() + i * 4096) as u64; - } - + let thread_ctxts = self.thread_ctxts.get_mut(); { - let regs = self.regs.read().unwrap(); + let regs = self.regs.read(); log::debug!("CAP_LOW: {:X}", regs.cap_low.read()); log::debug!("CAP_HIGH: {:X}", regs.cap_high.read()); log::debug!("VS: {:X}", regs.vs.read()); @@ -289,11 +297,11 @@ impl Nvme { } log::debug!("Disabling controller."); - self.regs.get_mut().unwrap().cc.writef(1, false); + self.regs.get_mut().cc.writef(1, false); log::trace!("Waiting for not ready."); loop { - let csts = self.regs.get_mut().unwrap().csts.read(); + let csts = self.regs.get_mut().csts.read(); log::trace!("CSTS: {:X}", csts); if csts & 1 == 1 { pause(); @@ -302,46 +310,41 @@ impl Nvme { } } - match self.interrupt_method.get_mut().unwrap() { + match self.interrupt_method.get_mut() { &mut InterruptMethod::Intx | InterruptMethod::Msi { .. } => { - self.regs.get_mut().unwrap().intms.write(0xFFFF_FFFF); - self.regs.get_mut().unwrap().intmc.write(0x0000_0001); + self.regs.get_mut().intms.write(0xFFFF_FFFF); + self.regs.get_mut().intmc.write(0x0000_0001); } &mut InterruptMethod::MsiX(ref mut cfg) => { cfg.table[0].unmask(); } } - for (qid, queue) in self.completion_queues.get_mut().unwrap().iter_mut() { - let &(ref cq, ref sq_ids) = &*queue.get_mut().unwrap(); - let data = &cq.data; - log::debug!( - "completion queue {}: {:X}, {}, (submission queue ids: {:?}", - qid, - data.physical(), - data.len(), - sq_ids - ); - } + for (qid, iv) in self.cq_ivs.get_mut().iter_mut() { + let ctxt = thread_ctxts.get(&0).unwrap().lock(); + let queues = ctxt.queues.borrow(); - for (qid, (queue, cq_id)) in self.submission_queues.get_mut().unwrap().iter_mut() { - let data = &queue.get_mut().unwrap().data; + let &(ref cq, ref sq) = queues.get(qid).unwrap(); log::debug!( - "submission queue {}: {:X}, {}, attached to CQID: {}", - qid, - data.physical(), - data.len(), - cq_id + "iv {iv} [cq {qid}: {:X}, {}] [sq {qid}: {:X}, {}]", + cq.data.physical(), + cq.data.len(), + sq.data.physical(), + sq.data.len() ); } { - let regs = self.regs.get_mut().unwrap(); - let submission_queues = self.submission_queues.get_mut().unwrap(); - let completion_queues = self.completion_queues.get_mut().unwrap(); + let main_ctxt = thread_ctxts.get(&0).unwrap().lock(); - let asq = submission_queues.get_mut(&0).unwrap().0.get_mut().unwrap(); - let (acq, _) = completion_queues.get_mut(&0).unwrap().get_mut().unwrap(); + for (i, prp) in main_ctxt.buffer_prp.borrow_mut().iter_mut().enumerate() { + *prp = (main_ctxt.buffer.borrow_mut().physical() + i * 4096) as u64; + } + + let regs = self.regs.get_mut(); + + let mut queues = main_ctxt.queues.borrow_mut(); + let (asq, acq) = queues.get_mut(&0).unwrap(); regs.aqa .write(((acq.data.len() as u32 - 1) << 16) | (asq.data.len() as u32 - 1)); regs.asq_low.write(asq.data.physical() as u32); @@ -359,11 +362,11 @@ impl Nvme { } log::debug!("Enabling controller."); - self.regs.get_mut().unwrap().cc.writef(1, true); + self.regs.get_mut().cc.writef(1, true); log::debug!("Waiting for ready"); loop { - let csts = self.regs.get_mut().unwrap().csts.read(); + let csts = self.regs.get_mut().csts.read(); log::debug!("CSTS: {:X}", csts); if csts & 1 == 0 { pause(); @@ -378,7 +381,7 @@ impl Nvme { /// # Panics /// Will panic if the same vector is called twice with different mask flags. pub fn set_vectors_masked(&self, vectors: impl IntoIterator) { - let mut interrupt_method_guard = self.interrupt_method.lock().unwrap(); + let mut interrupt_method_guard = self.interrupt_method.lock(); match &mut *interrupt_method_guard { &mut InterruptMethod::Intx => { @@ -394,9 +397,9 @@ impl Nvme { ); assert_eq!(vector, 0, "nvmed: internal error: nonzero vector on INTx#"); if mask { - self.regs.write().unwrap().intms.write(0x0000_0001); + self.regs.write().intms.write(0x0000_0001); } else { - self.regs.write().unwrap().intmc.write(0x0000_0001); + self.regs.write().intmc.write(0x0000_0001); } } &mut InterruptMethod::Msi { @@ -431,10 +434,10 @@ impl Nvme { } if to_mask != 0 { - self.regs.write().unwrap().intms.write(to_mask); + self.regs.write().intms.write(to_mask); } if to_clear != 0 { - self.regs.write().unwrap().intmc.write(to_clear); + self.regs.write().intmc.write(to_clear); } } &mut InterruptMethod::MsiX(ref mut cfg) => { @@ -451,226 +454,175 @@ impl Nvme { self.set_vectors_masked(std::iter::once((vector, masked))) } - #[cfg(not(feature = "async"))] - pub fn submit_and_complete_command NvmeCmd>( + pub async fn submit_and_complete_command( &self, sq_id: SqId, - cmd_init: F, + cmd_init: impl FnOnce(CmdId) -> NvmeCmd, ) -> NvmeComp { - // Submit command - let cmd = { - let sqs_read_guard = self.submission_queues.read().unwrap(); - let &(ref sq_lock, cq_id) = sqs_read_guard - .get(&sq_id) - .expect("nvmed: internal error: given SQ for SQ ID not there"); - let mut sq_guard = sq_lock.lock().unwrap(); - let sq = &mut *sq_guard; + NvmeExecutor::current().submit(sq_id, cmd_init(0)).await + } - assert!(!sq.is_full()); - - let cmd_id = u16::try_from(sq.tail) - .expect("nvmed: internal error: CQ has more than 2^16 entries"); - let cmd = cmd_init(cmd_id); - log::trace!( - "Sent submission queue entry (SQID {}): {:?} at {}", - sq_id, - cmd, - cmd_id - ); - let tail = sq.submit_unchecked(cmd); - let tail = u16::try_from(tail).unwrap(); - - // make sure that we register interest before the reactor can get notified - unsafe { self.submission_queue_tail(sq_id, tail) }; - - cmd - }; - - // Read completion - loop { - for (cq_id, completion_queue_lock) in self.completion_queues.read().unwrap().iter() { - if *cq_id != sq_id { - // Currently, CQ and SQ IDs have to match - continue; + pub async fn submit_and_complete_admin_command( + &self, + cmd_init: impl FnOnce(CmdId) -> NvmeCmd, + ) -> NvmeComp { + self.submit_and_complete_command(0, cmd_init).await + } + pub fn try_submit_raw( + &self, + ctxt: &ThreadCtxt, + sq_id: SqId, + cmd_init: impl FnOnce(CmdId) -> NvmeCmd, + fail: impl FnOnce(), + ) -> Option<(CqId, CmdId)> { + match ctxt.queues.borrow_mut().get_mut(&sq_id).unwrap() { + (sq, _cq) => { + if sq.is_full() { + fail(); + return None; } + let cmd_id = sq.tail; + let tail = sq.submit_unchecked(cmd_init(cmd_id)); - let mut completion_queue_guard = completion_queue_lock.lock().unwrap(); - let &mut (ref mut completion_queue, _) = &mut *completion_queue_guard; - - while let Some((head, entry)) = completion_queue.complete(Some((sq_id, cmd))) { - unsafe { self.completion_queue_head(*cq_id, head) }; - - log::trace!( - "Got completion queue entry (CQID {}): {:?} at {}", - cq_id, - entry, - head - ); - - assert_eq!(sq_id, { entry.sq_id }); - assert_eq!({ cmd.cid }, { entry.cid }); - - { - let submission_queues_read_lock = self.submission_queues.read().unwrap(); - // this lock is actually important, since it will block during submission from other - // threads. the lock won't be held for long by the submitters, but it still prevents - // the entry being lost before this reactor is actually able to respond: - let &(ref sq_lock, corresponding_cq_id) = submission_queues_read_lock.get(&{entry.sq_id}).expect("nvmed: internal error: queue returned from controller doesn't exist"); - assert_eq!(*cq_id, corresponding_cq_id); - let mut sq_guard = sq_lock.lock().unwrap(); - sq_guard.head = entry.sq_head; - } - - return entry; + // TODO: Submit in bulk + unsafe { + self.submission_queue_tail(sq_id, tail); } + Some((sq_id, cmd_id)) } - std::thread::yield_now(); } } - #[cfg(feature = "async")] - pub fn submit_and_complete_command NvmeCmd>( + pub async fn create_io_completion_queue( &self, - sq_id: SqId, - cmd_init: F, - ) -> NvmeComp { - use crate::nvme::cq_reactor::{CompletionFuture, CompletionFutureState}; - futures::executor::block_on(CompletionFuture { - state: CompletionFutureState::PendingSubmission { - cmd_init, - nvme: &self, - sq_id, - }, - }) - } + io_cq_id: CqId, + vector: Option, + ) -> NvmeCompQueue { + let queue = NvmeCompQueue::new().expect("nvmed: failed to allocate I/O completion queue"); - pub fn submit_and_complete_admin_command NvmeCmd>( - &self, - cmd_init: F, - ) -> NvmeComp { - self.submit_and_complete_command(0, cmd_init) - } - - pub fn create_io_completion_queue(&self, io_cq_id: CqId, vector: Option) { - let (ptr, len) = { - let mut completion_queues_guard = self.completion_queues.write().unwrap(); - - let queue_guard = completion_queues_guard - .entry(io_cq_id) - .or_insert_with(|| { - let queue = NvmeCompQueue::new() - .expect("nvmed: failed to allocate I/O completion queue"); - let sqs = SmallVec::new(); - Mutex::new((queue, sqs)) - }) - .get_mut() - .unwrap(); - - let &(ref queue, _) = &*queue_guard; - (queue.data.physical(), queue.data.len()) - }; - - let len = - u16::try_from(len).expect("nvmed: internal error: I/O CQ longer than 2^16 entries"); + let len = u16::try_from(queue.data.len()) + .expect("nvmed: internal error: I/O CQ longer than 2^16 entries"); let raw_len = len .checked_sub(1) .expect("nvmed: internal error: CQID 0 for I/O CQ"); - let comp = self.submit_and_complete_admin_command(|cid| { - NvmeCmd::create_io_completion_queue(cid, io_cq_id, ptr, raw_len, vector) - }); - - if let Some(vector) = vector { - self.cqs_for_ivs - .write() - .unwrap() - .entry(vector) - .or_insert_with(SmallVec::new) - .push(io_cq_id); - } - } - pub fn create_io_submission_queue(&self, io_sq_id: SqId, io_cq_id: CqId) { - let (ptr, len) = { - let mut submission_queues_guard = self.submission_queues.write().unwrap(); - - let (queue_lock, _) = submission_queues_guard.entry(io_sq_id).or_insert_with(|| { - ( - Mutex::new( - NvmeCmdQueue::new() - .expect("nvmed: failed to allocate I/O completion queue"), - ), + let comp = self + .submit_and_complete_admin_command(|cid| { + NvmeCmd::create_io_completion_queue( + cid, io_cq_id, + queue.data.physical(), + raw_len, + vector, ) - }); - let queue = queue_lock.get_mut().unwrap(); + }) + .await; - (queue.data.physical(), queue.data.len()) - }; + /*match comp.status.specific { + 1 => panic!("invalid queue identifier"), + 2 => panic!("invalid queue size"), + 8 => panic!("invalid interrupt vector"), + _ => (), + }*/ - let len = - u16::try_from(len).expect("nvmed: internal error: I/O SQ longer than 2^16 entries"); + queue + } + pub async fn create_io_submission_queue(&self, io_sq_id: SqId, io_cq_id: CqId) -> NvmeCmdQueue { + let q = NvmeCmdQueue::new().expect("failed to create submission queue"); + + let len = u16::try_from(q.data.len()) + .expect("nvmed: internal error: I/O SQ longer than 2^16 entries"); let raw_len = len .checked_sub(1) .expect("nvmed: internal error: SQID 0 for I/O SQ"); - let comp = self.submit_and_complete_admin_command(|cid| { - NvmeCmd::create_io_submission_queue(cid, io_sq_id, ptr, raw_len, io_cq_id) - }); + let comp = self + .submit_and_complete_admin_command(|cid| { + NvmeCmd::create_io_submission_queue( + cid, + io_sq_id, + q.data.physical(), + raw_len, + io_cq_id, + ) + }) + .await; + /*match comp.status.specific { + 0 => panic!("completion queue invalid"), + 1 => panic!("invalid queue identifier"), + 2 => panic!("invalid queue size"), + _ => (), + }*/ + + q } - pub fn init_with_queues(&self) -> BTreeMap { + pub async fn init_with_queues(&self) -> BTreeMap { log::trace!("preinit"); - self.identify_controller(); - let nsids = self.identify_namespace_list(0); + self.identify_controller().await; + + let nsids = self.identify_namespace_list(0).await; log::debug!("first commands"); let mut namespaces = BTreeMap::new(); for nsid in nsids.iter().copied() { - namespaces.insert(nsid, self.identify_namespace(nsid)); + namespaces.insert(nsid, self.identify_namespace(nsid).await); } // TODO: Multiple queues - self.create_io_completion_queue(1, Some(0)); - self.create_io_submission_queue(1, 1); + let cq = self.create_io_completion_queue(1, Some(0)).await; + log::trace!("created compq"); + let sq = self.create_io_submission_queue(1, 1).await; + log::trace!("created subq"); + self.thread_ctxts + .read() + .get(&0) + .unwrap() + .lock() + .queues + .borrow_mut() + .insert(1, (sq, cq)); + self.sq_ivs.write().insert(1, 0); + self.cq_ivs.write().insert(1, 0); namespaces } - fn namespace_rw( + async fn namespace_rw( &self, - namespace: NvmeNamespace, + ctxt: &ThreadCtxt, + namespace: &NvmeNamespace, lba: u64, blocks_1: u16, write: bool, ) -> Result<()> { let block_size = namespace.block_size; - let buffer_prp_guard = self.buffer_prp.lock().unwrap(); - + let prp = ctxt.buffer_prp.borrow_mut(); let bytes = ((blocks_1 as u64) + 1) * block_size; let (ptr0, ptr1) = if bytes <= 4096 { - (buffer_prp_guard[0], 0) + (prp[0], 0) } else if bytes <= 8192 { - (buffer_prp_guard[0], buffer_prp_guard[1]) + (prp[0], prp[1]) } else { - ( - buffer_prp_guard[0], - (buffer_prp_guard.physical() + 8) as u64, - ) + (prp[0], (prp.physical() + 8) as u64) }; let mut cmd = NvmeCmd::default(); - let comp = self.submit_and_complete_command(1, |cid| { - cmd = if write { - NvmeCmd::io_write(cid, namespace.id, lba, blocks_1, ptr0, ptr1) - } else { - NvmeCmd::io_read(cid, namespace.id, lba, blocks_1, ptr0, ptr1) - }; - cmd.clone() - }); + let comp = self + .submit_and_complete_command(1, |cid| { + cmd = if write { + NvmeCmd::io_write(cid, namespace.id, lba, blocks_1, ptr0, ptr1) + } else { + NvmeCmd::io_read(cid, namespace.id, lba, blocks_1, ptr0, ptr1) + }; + cmd.clone() + }) + .await; + let status = comp.status >> 1; if status == 0 { Ok(()) @@ -680,55 +632,59 @@ impl Nvme { } } - pub fn namespace_read( + pub async fn namespace_read( &self, - namespace: NvmeNamespace, + namespace: &NvmeNamespace, mut lba: u64, buf: &mut [u8], - ) -> Result> { + ) -> Result { + let ctxt = self.cur_thread_ctxt(); + let ctxt = ctxt.lock(); + let block_size = namespace.block_size as usize; - let buffer_guard = self.buffer.lock().unwrap(); - - for chunk in buf.chunks_mut(/*TODO: buffer_guard.len()*/ 8192) { + for chunk in buf.chunks_mut(/* TODO: buf len */ 8192) { let blocks = (chunk.len() + block_size - 1) / block_size; assert!(blocks > 0); assert!(blocks <= 0x1_0000); - self.namespace_rw(namespace, lba, (blocks - 1) as u16, false)?; + self.namespace_rw(&*ctxt, namespace, lba, (blocks - 1) as u16, false) + .await?; - chunk.copy_from_slice(&buffer_guard[..chunk.len()]); + chunk.copy_from_slice(&ctxt.buffer.borrow()[..chunk.len()]); lba += blocks as u64; } - Ok(Some(buf.len())) + Ok(buf.len()) } - pub fn namespace_write( + pub async fn namespace_write( &self, - namespace: NvmeNamespace, + namespace: &NvmeNamespace, mut lba: u64, buf: &[u8], - ) -> Result> { + ) -> Result { + let ctxt = self.cur_thread_ctxt(); + let ctxt = ctxt.lock(); + let block_size = namespace.block_size as usize; - let mut buffer_guard = self.buffer.lock().unwrap(); - - for chunk in buf.chunks(/*TODO: buffer_guard.len()*/ 8192) { + for chunk in buf.chunks(/* TODO: buf len */ 8192) { let blocks = (chunk.len() + block_size - 1) / block_size; assert!(blocks > 0); assert!(blocks <= 0x1_0000); - buffer_guard[..chunk.len()].copy_from_slice(chunk); + ctxt.buffer.borrow_mut()[..chunk.len()].copy_from_slice(chunk); - self.namespace_rw(namespace, lba, (blocks - 1) as u16, true)?; + self.namespace_rw(&*ctxt, namespace, lba, (blocks - 1) as u16, true) + .await?; lba += blocks as u64; } - Ok(Some(buf.len())) + Ok(buf.len()) } } diff --git a/storage/nvmed/src/nvme/queues.rs b/storage/nvmed/src/nvme/queues.rs index 133c42effe..9d9cf2e2cb 100644 --- a/storage/nvmed/src/nvme/queues.rs +++ b/storage/nvmed/src/nvme/queues.rs @@ -1,3 +1,4 @@ +use std::cell::UnsafeCell; use std::ptr; use syscall::Result; @@ -49,7 +50,7 @@ pub struct NvmeComp { /// Completion queue pub struct NvmeCompQueue { - pub data: Dma<[NvmeComp]>, + pub data: Dma<[UnsafeCell]>, pub head: u16, pub phase: bool, } @@ -64,15 +65,8 @@ impl NvmeCompQueue { } /// Get a new completion queue entry, or return None if no entry is available yet. - pub(crate) fn complete(&mut self, cmd_opt: Option<(u16, NvmeCmd)>) -> Option<(u16, NvmeComp)> { - let entry = unsafe { ptr::read_volatile(self.data.as_ptr().add(self.head as usize)) }; - - //HACK FOR SOMETIMES RETURNING INVALID DATA ON QEMU! - if let Some((sq_id, cmd)) = cmd_opt { - if entry.sq_id != sq_id || entry.cid != cmd.cid { - return None; - } - } + pub(crate) fn complete(&mut self) -> Option<(u16, NvmeComp)> { + let entry = unsafe { ptr::read_volatile(self.data[usize::from(self.head)].get()) }; if ((entry.status & 1) == 1) == self.phase { self.head = (self.head + 1) % (self.data.len() as u16); @@ -86,10 +80,10 @@ impl NvmeCompQueue { } /// Get a new CQ entry, busy waiting until an entry appears. - fn complete_spin(&mut self, cmd_opt: Option<(u16, NvmeCmd)>) -> (u16, NvmeComp) { + pub fn complete_spin(&mut self) -> (u16, NvmeComp) { log::debug!("Waiting for new CQ entry"); loop { - if let Some(some) = self.complete(cmd_opt) { + if let Some(some) = self.complete() { return some; } else { unsafe { @@ -102,7 +96,7 @@ impl NvmeCompQueue { /// Submission queue pub struct NvmeCmdQueue { - pub data: Dma<[NvmeCmd]>, + pub data: Dma<[UnsafeCell]>, pub tail: u16, pub head: u16, } @@ -126,8 +120,32 @@ impl NvmeCmdQueue { /// Add a new submission command entry to the queue. The caller must ensure that the queue have free /// entries; this can be checked using `is_full`. pub fn submit_unchecked(&mut self, entry: NvmeCmd) -> u16 { - unsafe { ptr::write_volatile(&mut self.data[self.tail as usize] as *mut _, entry) } + unsafe { ptr::write_volatile(self.data[usize::from(self.tail)].get(), entry) } self.tail = (self.tail + 1) % (self.data.len() as u16); self.tail } } + +#[derive(Debug)] +pub enum Status { + GenericCmdStatus(u8), + CommandSpecificStatus(u8), + IntegrityError(u8), + PathRelatedStatus(u8), + Rsvd(u8), + Vendor(u8), +} +impl Status { + pub fn parse(raw: u16) -> Self { + let code = (raw >> 1) as u8; + match (raw >> 9) & 0b111 { + 0 => Self::GenericCmdStatus(code), + 1 => Self::CommandSpecificStatus(code), + 2 => Self::IntegrityError(code), + 3 => Self::PathRelatedStatus(code), + 4..=6 => Self::Rsvd(code), + 7 => Self::Vendor(code), + _ => unreachable!(), + } + } +} diff --git a/storage/usbscsid/src/main.rs b/storage/usbscsid/src/main.rs index 807ab11c38..6da8c36fe6 100644 --- a/storage/usbscsid/src/main.rs +++ b/storage/usbscsid/src/main.rs @@ -1,7 +1,7 @@ use std::collections::BTreeMap; use std::env; -use driver_block::{Disk, DiskScheme}; +use driver_block::{Disk, DiskScheme, ExecutorTrait}; use syscall::{Error, EIO}; use xhcid_interface::{ConfigureEndpointsReq, PortId, XhciClientHandle}; @@ -106,6 +106,7 @@ fn daemon(daemon: redox_daemon::Daemon, scheme: String, port: PortId, protocol: protocol: &mut *protocol, }, )]), + &driver_block::FuturesExecutor, ); //libredox::call::setrens(0, 0).expect("nvmed: failed to enter null namespace"); @@ -120,7 +121,9 @@ fn daemon(daemon: redox_daemon::Daemon, scheme: String, port: PortId, protocol: for event in event_queue { match event.unwrap().user_data { - Event::Scheme => scheme.tick().unwrap(), + Event::Scheme => driver_block::FuturesExecutor + .block_on(scheme.tick()) + .unwrap(), } } @@ -141,9 +144,9 @@ impl Disk for UsbDisk<'_> { self.scsi.get_disk_size() } - fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result> { + async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result { match self.scsi.read(self.protocol, block, buffer) { - Ok(bytes_read) => Ok(Some(bytes_read as usize)), + Ok(bytes_read) => Ok(bytes_read as usize), Err(err) => { eprintln!("usbscsid: READ IO ERROR: {err}"); Err(Error::new(EIO)) @@ -151,9 +154,9 @@ impl Disk for UsbDisk<'_> { } } - fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result> { + async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result { match self.scsi.write(self.protocol, block, buffer) { - Ok(bytes_written) => Ok(Some(bytes_written as usize)), + Ok(bytes_written) => Ok(bytes_written as usize), Err(err) => { eprintln!("usbscsid: WRITE IO ERROR: {err}"); Err(Error::new(EIO)) diff --git a/storage/virtio-blkd/src/main.rs b/storage/virtio-blkd/src/main.rs index 7503b477d7..88fbb52bdf 100644 --- a/storage/virtio-blkd/src/main.rs +++ b/storage/virtio-blkd/src/main.rs @@ -154,6 +154,7 @@ fn deamon(deamon: redox_daemon::Daemon) -> anyhow::Result<()> { let mut scheme = DiskScheme::new( scheme_name, BTreeMap::from([(0, VirtioDisk::new(queue, device_space))]), + &driver_block::FuturesExecutor, ); libredox::call::setrens(0, 0).expect("nvmed: failed to enter null namespace"); @@ -170,7 +171,7 @@ fn deamon(deamon: redox_daemon::Daemon) -> anyhow::Result<()> { for event in event_queue { match event.unwrap().user_data { - Event::Scheme => scheme.tick().unwrap(), + Event::Scheme => futures::executor::block_on(scheme.tick()).unwrap(), } } diff --git a/storage/virtio-blkd/src/scheme.rs b/storage/virtio-blkd/src/scheme.rs index 074f2ff61c..ec4ecf732d 100644 --- a/storage/virtio-blkd/src/scheme.rs +++ b/storage/virtio-blkd/src/scheme.rs @@ -93,15 +93,11 @@ impl driver_block::Disk for VirtioDisk<'_> { self.cfg.capacity() * u64::from(self.cfg.block_size()) } - fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result> { - Ok(Some(futures::executor::block_on( - self.queue.read(block, buffer), - ))) + async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result { + Ok(self.queue.read(block, buffer).await) } - fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result> { - Ok(Some(futures::executor::block_on( - self.queue.write(block, buffer), - ))) + async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result { + Ok(self.queue.write(block, buffer).await) } }