Merge branch 'nvme-executor' into 'master'

Integrate nvmed executor with reactor, moving towards thread-per-core

See merge request redox-os/drivers!155
This commit is contained in:
Jeremy Soller
2025-03-29 14:03:50 +00:00
22 changed files with 1280 additions and 648 deletions
Generated
+180 -127
View File
@@ -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",
]
+3 -1
View File
@@ -1,7 +1,9 @@
[workspace]
members = [
"acpid",
"common",
"executor",
"acpid",
"hwd",
"pcid",
"pcid-spawner",
+12
View File
@@ -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"
+396
View File
@@ -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<LocalExecutor<Self>>;
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<Hw: Hardware> {
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<HashMap<Hw::SqId, VecDeque<FutIdx>>>,
awaiting_completion:
RefCell<HashMap<Hw::CqId, HashMap<Hw::CmdId, (FutIdx, NonNull<Option<Hw::Cqe>>)>>>,
external_event: RefCell<HashMap<EventUserData, (FutIdx, NonNull<EventFlags>)>>,
next_user_data: Cell<usize>,
ready_futures: RefCell<VecDeque<FutIdx>>,
futures: RefCell<Slab<Pin<Box<dyn Future<Output = ()> + 'static>>>>,
is_polling: Cell<bool>,
}
impl<Hw: Hardware> LocalExecutor<Hw> {
pub fn register_external_event(
&self,
fd: usize,
flags: event::EventFlags,
) -> ExternalEventSource<Hw> {
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<Self> {
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::<Hw>(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<Output = ()> + '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<Output = O> + 'a) -> O {
let retval = Rc::new(RefCell::new(None));
let retval2 = Rc::clone(&retval);
let idx = self.futures.borrow_mut().insert({
let t1: Pin<Box<dyn Future<Output = ()> + '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<Box<dyn Future<Output = ()> + '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::<usize>()];
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::<Hw> {
state: State::<Hw>::Submitting { sq_id, cmd },
comp: None,
_not_send: PhantomData,
}
.await
}
}
struct CqeFuture<Hw: Hardware> {
pub state: State<Hw>,
pub comp: Option<Hw::Cqe>,
pub _not_send: PhantomData<*const ()>,
}
enum State<Hw: Hardware> {
Submitting { sq_id: Hw::SqId, cmd: Hw::Sqe },
Completing { cq_id: Hw::CqId, cmd_id: Hw::CmdId },
}
fn current_executor_and_idx<Hw: Hardware>(
cx: &mut task::Context<'_>,
) -> (Rc<LocalExecutor<Hw>>, 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<Hw: Hardware> Future for CqeFuture<Hw> {
type Output = Hw::Cqe;
fn poll(self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> task::Poll<Self::Output> {
let this = unsafe { self.get_unchecked_mut() };
let (executor, idx) = current_executor_and_idx::<Hw>(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<Hw: Hardware>(idx: *const ()) -> task::RawWaker {
task::RawWaker::new(idx, Hw::vtable())
}
unsafe fn vt_drop(_idx: *const ()) {}
unsafe fn vt_wake<Hw: Hardware>(idx: *const ()) {
Hw::current()
.ready_futures
.borrow_mut()
.push_back(idx as FutIdx);
}
fn waker<Hw: Hardware>(idx: FutIdx) -> task::Waker {
unsafe { task::Waker::from_raw(task::RawWaker::new(idx as *const (), Hw::vtable())) }
}
pub const fn vtable<Hw: Hardware>() -> task::RawWakerVTable {
task::RawWakerVTable::new(vt_clone::<Hw>, vt_wake::<Hw>, vt_wake::<Hw>, vt_drop)
}
pub struct ExternalEventSource<Hw: Hardware> {
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<Hw: Hardware> ExternalEventSource<Hw> {
fn poll_next(self: Pin<&mut Self>, cx: &mut task::Context) -> task::Poll<Option<Event>> {
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::<Hw>(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<Event> {
core::future::poll_fn(|cx| self.as_mut().poll_next(cx)).await
}
}
pub fn init_raw<Hw: Hardware>(
global_ctxt: Hw::GlobalCtxt,
vector: Hw::Iv,
intx: bool,
irq_handle: File,
) -> LocalExecutor<Hw> {
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),
}
}
+5 -4
View File
@@ -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<Option<usize>> {
async fn read(&mut self, block: u64, buffer: &mut [u8]) -> Result<usize> {
//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<Option<usize>> {
async fn write(&mut self, block: u64, buffer: &[u8]) -> Result<usize> {
//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(),
}
}
+3 -3
View File
@@ -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<Option<usize>> {
async fn read(&mut self, block: u64, buffer: &mut [u8]) -> Result<usize> {
// 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<Option<usize>> {
async fn write(&mut self, _block: u64, _buffer: &[u8]) -> Result<usize> {
Err(Error::new(EBADF)) // TODO: Implement writing
}
}
+36 -5
View File
@@ -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<Box<dyn Disk>>) {
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<usize> {
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<usize> {
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<AnyDisk>) {
let hba_mem = unsafe { &mut *(base as *mut HbaMem) };
hba_mem.init();
let pi = hba_mem.pi.read();
let disks: Vec<Box<dyn Disk>> = (0..hba_mem.ports.len())
let disks: Vec<AnyDisk> = (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<Box<dyn Disk>> = match port_type {
let disk: Option<AnyDisk> = 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
+4 -6
View File
@@ -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 {
+8 -2
View File
@@ -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"
+160 -124
View File
@@ -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<Option<usize>>;
fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result<Option<usize>>;
async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result<usize>;
async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result<usize>;
}
impl<T: Disk + ?Sized> Disk for Box<T> {
@@ -84,12 +87,12 @@ impl<T: Disk + ?Sized> Disk for Box<T> {
(**self).size()
}
fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result<Option<usize>> {
(**self).read(block, buffer)
async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result<usize> {
(**self).read(block, buffer).await
}
fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result<Option<usize>> {
(**self).write(block, buffer)
async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result<usize> {
(**self).write(block, buffer).await
}
}
@@ -99,18 +102,19 @@ pub struct DiskWrapper<T> {
}
impl<T: Disk> DiskWrapper<T> {
pub fn pt(disk: &mut T) -> Option<PartitionTable> {
pub fn pt(disk: &mut T, executor: &impl ExecutorTrait) -> Option<PartitionTable> {
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<u64> {
let size = i64::try_from(self.disk.size()).or(Err(io::Error::new(
io::ErrorKind::Other,
@@ -129,7 +133,7 @@ impl<T: Disk> DiskWrapper<T> {
}
}
// 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<usize> {
let blksize = self.disk.block_size();
let size_in_blocks = self.disk.size() / u64::from(blksize);
@@ -141,17 +145,9 @@ impl<T: Disk> DiskWrapper<T> {
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<T: Disk> DiskWrapper<T> {
}
}
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<T: Disk> DiskWrapper<T> {
self.disk.size()
}
pub fn read(
pub async fn read(
&mut self,
part_num: Option<usize>,
block: u64,
buf: &mut [u8],
) -> syscall::Result<Option<usize>> {
) -> syscall::Result<usize> {
if buf.len() as u64 % u64::from(self.disk.block_size()) != 0 {
return Err(Error::new(EINVAL));
}
@@ -214,18 +217,18 @@ impl<T: Disk> DiskWrapper<T> {
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<usize>,
block: u64,
buf: &[u8],
) -> syscall::Result<Option<usize>> {
) -> syscall::Result<usize> {
if buf.len() as u64 % u64::from(self.disk.block_size()) != 0 {
return Err(Error::new(EINVAL));
}
@@ -245,9 +248,9 @@ impl<T: Disk> DiskWrapper<T> {
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<T> {
disks: BTreeMap<u32, DiskWrapper<T>>,
handles: BTreeMap<usize, Handle>,
next_id: usize,
blocked: Vec<CallRequest>,
}
pub trait ExecutorTrait {
fn block_on<'a, O: 'a>(&self, fut: impl IntoFuture<Output = O> + 'a) -> O;
}
impl<Hw: executor::Hardware> ExecutorTrait for LocalExecutor<Hw> {
fn block_on<'a, O: 'a>(&self, fut: impl IntoFuture<Output = O> + '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<Output = O> + '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<Output = O> + '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<T: Disk> DiskScheme<T> {
pub fn new(scheme_name: String, disks: BTreeMap<u32, T>) -> Self {
pub fn new(
scheme_name: String,
disks: BTreeMap<u32, T>,
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<T: Disk> DiskScheme<T> {
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<T: Disk> DiskScheme<T> {
/// 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<T: Disk> DiskScheme<T> {
}
}
impl<T: Disk> SchemeBlock for DiskScheme<T> {
fn xopen(
&mut self,
path_str: &str,
flags: usize,
ctx: &CallerCtx,
) -> Result<Option<OpenResult>> {
impl<T: Disk> SchemeAsync for DiskScheme<T> {
async fn open(&mut self, path_str: &str, flags: usize, ctx: &CallerCtx) -> Result<OpenResult> {
if ctx.uid != 0 {
return Err(Error::new(EACCES));
}
@@ -446,18 +472,27 @@ impl<T: Disk> SchemeBlock for DiskScheme<T> {
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<DirentBuf<&'buf mut [u8]>> {
// TODO
Err(Error::new(EOPNOTSUPP))
}
fn fstat(&mut self, id: usize, stat: &mut Stat) -> Result<Option<usize>> {
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<T: Disk> SchemeBlock for DiskScheme<T> {
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<T: Disk> SchemeBlock for DiskScheme<T> {
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<Option<usize>> {
async fn fpath(&mut self, id: usize, buf: &mut [u8], _ctx: &CallerCtx) -> Result<usize> {
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<T: Disk> SchemeBlock for DiskScheme<T> {
}
}
Ok(Some(i))
Ok(i)
}
fn read(
async fn read(
&mut self,
id: usize,
buf: &mut [u8],
offset: u64,
_fcntl_flags: u32,
) -> Result<Option<usize>> {
_ctx: &CallerCtx,
) -> Result<usize> {
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<T: Disk> SchemeBlock for DiskScheme<T> {
.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<Option<usize>> {
_ctx: &CallerCtx,
) -> Result<usize> {
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<Option<u64>> {
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<u64> {
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<T: Disk> DiskScheme<T> {
fn on_close(&mut self, id: usize) {
self.handles.remove(&id);
impl<D: Disk> DiskScheme<D> {
pub fn on_close(&mut self, id: usize) {
let _ = self.handles.remove(&id);
}
}
+6 -4
View File
@@ -177,7 +177,8 @@ impl Disk for AtaDisk {
self.size
}
fn read(&mut self, start_block: u64, buffer: &mut [u8]) -> Result<Option<usize>> {
// NOTE: not async
async fn read(&mut self, start_block: u64, buffer: &mut [u8]) -> Result<usize> {
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<Option<usize>> {
// NOTE: not async
async fn write(&mut self, start_block: u64, buffer: &[u8]) -> Result<usize> {
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)
}
}
+30 -6
View File
@@ -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<Box<dyn Disk>> = 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<usize> {
let AnyDisk::Ata(a) = self;
a.write(block, buffer).await
}
async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result<usize> {
let AnyDisk::Ata(a) = self;
a.read(block, buffer).await
}
}
let mut disks: Vec<AnyDisk> = 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);
+7 -5
View File
@@ -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<Option<usize>> {
async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result<usize> {
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<Option<usize>> {
async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result<usize> {
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(),
}
}
+8 -8
View File
@@ -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 = []
+54 -47
View File
@@ -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<Nvme>, NvmeNamespace);
struct NvmeDisk {
nvme: Arc<Nvme>,
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<Option<usize>> {
self.0.namespace_read(self.1, block, buffer)
async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result<usize> {
self.nvme.namespace_read(&self.ns, block, buffer).await
}
fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result<Option<usize>> {
self.0.namespace_write(self.1, block, buffer)
async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result<usize> {
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);
}
+82
View File
@@ -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<Nvme>;
fn mask_vector(ctxt: &Arc<Nvme>, iv: Self::Iv) {
ctxt.set_vector_masked(iv, true)
}
fn unmask_vector(ctxt: &Arc<Nvme>, 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<executor::LocalExecutor<Self>> {
THE_EXECUTOR.with(|exec| Rc::clone(exec.borrow().as_ref().unwrap()))
}
fn try_submit(
nvme: &Arc<Nvme>,
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<Nvme>, 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<Nvme>, id: Self::CqId) -> Self::SqId {
id
}
}
static VTABLE: std::task::RawWakerVTable = executor::vtable::<NvmeHw>();
thread_local! {
static THE_EXECUTOR: RefCell<Option<Rc<LocalExecutor<NvmeHw>>>> = RefCell::new(None);
}
pub type NvmeExecutor = LocalExecutor<NvmeHw>;
pub fn init(nvme: Arc<Nvme>, iv: u16, intx: bool, irq_handle: File) -> Rc<LocalExecutor<NvmeHw>> {
let this = Rc::new(executor::init_raw(nvme, iv, intx, irq_handle));
THE_EXECUTOR.with(|exec| *exec.borrow_mut() = Some(Rc::clone(&this)));
this
}
+20 -14
View File
@@ -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<IdentifyControllerData> = 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<u32> {
pub async fn identify_namespace_list(&self, base: u32) -> Vec<u32> {
// 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<IdentifyNamespaceData> = 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();
+219 -263
View File
@@ -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<u16, File>),
@@ -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<InterruptMethod>,
pcid_interface: Mutex<PciFunctionHandle>,
regs: RwLock<&'static mut NvmeRegs>,
pub(crate) submission_queues: RwLock<BTreeMap<SqId, (Mutex<NvmeCmdQueue>, CqId)>>,
pub(crate) completion_queues:
RwLock<BTreeMap<CqId, Mutex<(NvmeCompQueue, SmallVec<[SqId; 16]>)>>>,
sq_ivs: RwLock<HashMap<SqId, Iv>>,
cq_ivs: RwLock<HashMap<CqId, Iv>>,
// maps interrupt vectors with the completion queues they have
cqs_for_ivs: RwLock<BTreeMap<u16, SmallVec<[CqId; 4]>>>,
buffer: Mutex<Dma<[u8; 512 * 4096]>>, // 2MB of buffer
buffer_prp: Mutex<Dma<[u64; 512]>>, // 4KB of PRP for the buffer
reactor_sender: Sender<cq_reactor::NotifReq>,
thread_ctxts: RwLock<HashMap<Iv, Arc<ReentrantMutex<ThreadCtxt>>>>,
next_sqid: AtomicSqId,
next_cqid: AtomicCqId,
next_avail_submission_epoch: AtomicU64,
}
pub struct ThreadCtxt {
buffer: RefCell<Dma<[u8; 512 * 4096]>>, // 2MB of buffer
buffer_prp: RefCell<Dma<[u64; 512]>>, // 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<HashMap<u16, (NvmeCmdQueue, NvmeCompQueue)>>,
}
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<NotifReq>,
) -> Result<Self> {
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<u32>)).write(value);
}
fn cur_thread_ctxt(&self) -> Arc<ReentrantMutex<ThreadCtxt>> {
// 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<Item = (u16, bool)>) {
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<F: FnOnce(CmdId) -> 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<F: FnOnce(CmdId) -> 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<Iv>,
) -> NvmeCompQueue {
let queue = NvmeCompQueue::new().expect("nvmed: failed to allocate I/O completion queue");
pub fn submit_and_complete_admin_command<F: FnOnce(CmdId) -> 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<u16>) {
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<u32, NvmeNamespace> {
pub async fn init_with_queues(&self) -> BTreeMap<u32, NvmeNamespace> {
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<Option<usize>> {
) -> Result<usize> {
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<Option<usize>> {
) -> Result<usize> {
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())
}
}
+32 -14
View File
@@ -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<NvmeComp>]>,
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<NvmeCmd>]>,
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!(),
}
}
}
+9 -6
View File
@@ -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<Option<usize>> {
async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result<usize> {
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<Option<usize>> {
async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result<usize> {
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))
+2 -1
View File
@@ -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(),
}
}
+4 -8
View File
@@ -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<Option<usize>> {
Ok(Some(futures::executor::block_on(
self.queue.read(block, buffer),
)))
async fn read(&mut self, block: u64, buffer: &mut [u8]) -> syscall::Result<usize> {
Ok(self.queue.read(block, buffer).await)
}
fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result<Option<usize>> {
Ok(Some(futures::executor::block_on(
self.queue.write(block, buffer),
)))
async fn write(&mut self, block: u64, buffer: &[u8]) -> syscall::Result<usize> {
Ok(self.queue.write(block, buffer).await)
}
}