diff --git a/Cargo.lock b/Cargo.lock index 1c1fd29..0ed0627 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -58,6 +58,23 @@ dependencies = [ "libc 0.2.42 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "base64" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "byteorder 1.2.3 (registry+https://github.com/rust-lang/crates.io-index)", +] + +[[package]] +name = "base64" +version = "0.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "byteorder 1.2.3 (registry+https://github.com/rust-lang/crates.io-index)", + "safemem 0.2.0 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "base64" version = "0.9.2" @@ -67,6 +84,11 @@ dependencies = [ "safemem 0.2.0 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "bitflags" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" + [[package]] name = "bitflags" version = "1.0.3" @@ -159,6 +181,23 @@ dependencies = [ "vec_map 0.8.1 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "core-foundation" +version = "0.2.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "core-foundation-sys 0.2.3 (registry+https://github.com/rust-lang/crates.io-index)", + "libc 0.2.42 (registry+https://github.com/rust-lang/crates.io-index)", +] + +[[package]] +name = "core-foundation-sys" +version = "0.2.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "libc 0.2.42 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "crossbeam-deque" version = "0.3.1" @@ -234,6 +273,19 @@ dependencies = [ "redox_syscall 0.1.40 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "foreign-types" +version = "0.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "foreign-types-shared 0.1.1 (registry+https://github.com/rust-lang/crates.io-index)", +] + +[[package]] +name = "foreign-types-shared" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" + [[package]] name = "fs_extra" version = "1.1.0" @@ -280,6 +332,24 @@ dependencies = [ "quick-error 1.2.2 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "hyper" +version = "0.10.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "base64 0.6.0 (registry+https://github.com/rust-lang/crates.io-index)", + "httparse 1.2.4 (registry+https://github.com/rust-lang/crates.io-index)", + "language-tags 0.2.2 (registry+https://github.com/rust-lang/crates.io-index)", + "log 0.3.9 (registry+https://github.com/rust-lang/crates.io-index)", + "mime 0.2.6 (registry+https://github.com/rust-lang/crates.io-index)", + "num_cpus 1.8.0 (registry+https://github.com/rust-lang/crates.io-index)", + "time 0.1.40 (registry+https://github.com/rust-lang/crates.io-index)", + "traitobject 0.1.0 (registry+https://github.com/rust-lang/crates.io-index)", + "typeable 0.1.2 (registry+https://github.com/rust-lang/crates.io-index)", + "unicase 1.4.2 (registry+https://github.com/rust-lang/crates.io-index)", + "url 1.7.1 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "hyper" version = "0.11.27" @@ -306,6 +376,16 @@ dependencies = [ "want 0.0.4 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "idna" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "matches 0.1.7 (registry+https://github.com/rust-lang/crates.io-index)", + "unicode-bidi 0.3.4 (registry+https://github.com/rust-lang/crates.io-index)", + "unicode-normalization 0.1.7 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "iovec" version = "0.1.2" @@ -334,6 +414,11 @@ name = "language-tags" version = "0.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" +[[package]] +name = "lazy_static" +version = "0.2.11" +source = "registry+https://github.com/rust-lang/crates.io-index" + [[package]] name = "lazy_static" version = "1.0.1" @@ -387,6 +472,11 @@ dependencies = [ "linked-hash-map 0.4.2 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "matches" +version = "0.1.7" +source = "registry+https://github.com/rust-lang/crates.io-index" + [[package]] name = "memchr" version = "2.0.1" @@ -409,6 +499,14 @@ name = "memoffset" version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" +[[package]] +name = "mime" +version = "0.2.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "log 0.3.9 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "mime" version = "0.3.7" @@ -476,6 +574,20 @@ dependencies = [ "winapi 0.3.5 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "native-tls" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "lazy_static 0.2.11 (registry+https://github.com/rust-lang/crates.io-index)", + "libc 0.2.42 (registry+https://github.com/rust-lang/crates.io-index)", + "openssl 0.9.24 (registry+https://github.com/rust-lang/crates.io-index)", + "schannel 0.1.13 (registry+https://github.com/rust-lang/crates.io-index)", + "security-framework 0.1.16 (registry+https://github.com/rust-lang/crates.io-index)", + "security-framework-sys 0.1.16 (registry+https://github.com/rust-lang/crates.io-index)", + "tempdir 0.3.7 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "net2" version = "0.2.32" @@ -524,6 +636,29 @@ dependencies = [ "libc 0.2.42 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "openssl" +version = "0.9.24" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "bitflags 0.9.1 (registry+https://github.com/rust-lang/crates.io-index)", + "foreign-types 0.3.2 (registry+https://github.com/rust-lang/crates.io-index)", + "lazy_static 1.0.1 (registry+https://github.com/rust-lang/crates.io-index)", + "libc 0.2.42 (registry+https://github.com/rust-lang/crates.io-index)", + "openssl-sys 0.9.35 (registry+https://github.com/rust-lang/crates.io-index)", +] + +[[package]] +name = "openssl-sys" +version = "0.9.35" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "cc 1.0.17 (registry+https://github.com/rust-lang/crates.io-index)", + "libc 0.2.42 (registry+https://github.com/rust-lang/crates.io-index)", + "pkg-config 0.3.11 (registry+https://github.com/rust-lang/crates.io-index)", + "vcpkg 0.2.3 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "percent-encoding" version = "1.0.1" @@ -578,6 +713,7 @@ dependencies = [ "serde_json 1.0.20 (registry+https://github.com/rust-lang/crates.io-index)", "tokio-core 0.1.17 (registry+https://github.com/rust-lang/crates.io-index)", "tokio-timer 0.2.4 (registry+https://github.com/rust-lang/crates.io-index)", + "websocket 0.20.3 (registry+https://github.com/rust-lang/crates.io-index)", ] [[package]] @@ -618,6 +754,7 @@ dependencies = [ "tokio-uds 0.1.7 (registry+https://github.com/rust-lang/crates.io-index)", "toml 0.4.6 (registry+https://github.com/rust-lang/crates.io-index)", "walkdir 2.1.4 (registry+https://github.com/rust-lang/crates.io-index)", + "websocket 0.20.3 (registry+https://github.com/rust-lang/crates.io-index)", ] [[package]] @@ -749,6 +886,15 @@ dependencies = [ "winapi 0.3.5 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "schannel" +version = "0.1.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "lazy_static 1.0.1 (registry+https://github.com/rust-lang/crates.io-index)", + "winapi 0.3.5 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "scoped-tls" version = "0.1.2" @@ -759,6 +905,26 @@ name = "scopeguard" version = "0.3.3" source = "registry+https://github.com/rust-lang/crates.io-index" +[[package]] +name = "security-framework" +version = "0.1.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "core-foundation 0.2.3 (registry+https://github.com/rust-lang/crates.io-index)", + "core-foundation-sys 0.2.3 (registry+https://github.com/rust-lang/crates.io-index)", + "libc 0.2.42 (registry+https://github.com/rust-lang/crates.io-index)", + "security-framework-sys 0.1.16 (registry+https://github.com/rust-lang/crates.io-index)", +] + +[[package]] +name = "security-framework-sys" +version = "0.1.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "core-foundation-sys 0.2.3 (registry+https://github.com/rust-lang/crates.io-index)", + "libc 0.2.42 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "serde" version = "1.0.66" @@ -801,6 +967,11 @@ dependencies = [ "serde 1.0.66 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "sha1" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" + [[package]] name = "slab" version = "0.3.0" @@ -1097,6 +1268,17 @@ dependencies = [ "tokio-executor 0.1.2 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "tokio-tls" +version = "0.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "futures 0.1.21 (registry+https://github.com/rust-lang/crates.io-index)", + "native-tls 0.1.5 (registry+https://github.com/rust-lang/crates.io-index)", + "tokio-core 0.1.17 (registry+https://github.com/rust-lang/crates.io-index)", + "tokio-io 0.1.6 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "tokio-udp" version = "0.1.0" @@ -1134,16 +1316,34 @@ dependencies = [ "serde 1.0.66 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "traitobject" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" + [[package]] name = "try-lock" version = "0.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" +[[package]] +name = "typeable" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" + [[package]] name = "ucd-util" version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" +[[package]] +name = "unicase" +version = "1.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "version_check 0.1.3 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "unicase" version = "2.1.0" @@ -1152,6 +1352,19 @@ dependencies = [ "version_check 0.1.3 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "unicode-bidi" +version = "0.3.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "matches 0.1.7 (registry+https://github.com/rust-lang/crates.io-index)", +] + +[[package]] +name = "unicode-normalization" +version = "0.1.7" +source = "registry+https://github.com/rust-lang/crates.io-index" + [[package]] name = "unicode-width" version = "0.1.5" @@ -1170,6 +1383,16 @@ dependencies = [ "void 1.0.2 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "url" +version = "1.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "idna 0.1.5 (registry+https://github.com/rust-lang/crates.io-index)", + "matches 0.1.7 (registry+https://github.com/rust-lang/crates.io-index)", + "percent-encoding 1.0.1 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "utf8-ranges" version = "1.0.0" @@ -1214,6 +1437,27 @@ dependencies = [ "try-lock 0.1.0 (registry+https://github.com/rust-lang/crates.io-index)", ] +[[package]] +name = "websocket" +version = "0.20.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +dependencies = [ + "base64 0.5.2 (registry+https://github.com/rust-lang/crates.io-index)", + "bitflags 0.9.1 (registry+https://github.com/rust-lang/crates.io-index)", + "byteorder 1.2.3 (registry+https://github.com/rust-lang/crates.io-index)", + "bytes 0.4.8 (registry+https://github.com/rust-lang/crates.io-index)", + "futures 0.1.21 (registry+https://github.com/rust-lang/crates.io-index)", + "hyper 0.10.13 (registry+https://github.com/rust-lang/crates.io-index)", + "native-tls 0.1.5 (registry+https://github.com/rust-lang/crates.io-index)", + "rand 0.3.22 (registry+https://github.com/rust-lang/crates.io-index)", + "sha1 0.2.0 (registry+https://github.com/rust-lang/crates.io-index)", + "tokio-core 0.1.17 (registry+https://github.com/rust-lang/crates.io-index)", + "tokio-io 0.1.6 (registry+https://github.com/rust-lang/crates.io-index)", + "tokio-tls 0.1.4 (registry+https://github.com/rust-lang/crates.io-index)", + "unicase 1.4.2 (registry+https://github.com/rust-lang/crates.io-index)", + "url 1.7.1 (registry+https://github.com/rust-lang/crates.io-index)", +] + [[package]] name = "winapi" version = "0.2.8" @@ -1276,7 +1520,10 @@ dependencies = [ "checksum atty 0.2.10 (registry+https://github.com/rust-lang/crates.io-index)" = "2fc4a1aa4c24c0718a250f0681885c1af91419d242f29eb8f2ab28502d80dbd1" "checksum backtrace 0.3.8 (registry+https://github.com/rust-lang/crates.io-index)" = "dbdd17cd962b570302f5297aea8648d5923e22e555c2ed2d8b2e34eca646bf6d" "checksum backtrace-sys 0.1.23 (registry+https://github.com/rust-lang/crates.io-index)" = "bff67d0c06556c0b8e6b5f090f0eac52d950d9dfd1d35ba04e4ca3543eaf6a7e" +"checksum base64 0.5.2 (registry+https://github.com/rust-lang/crates.io-index)" = "30e93c03064e7590d0466209155251b90c22e37fab1daf2771582598b5827557" +"checksum base64 0.6.0 (registry+https://github.com/rust-lang/crates.io-index)" = "96434f987501f0ed4eb336a411e0631ecd1afa11574fe148587adc4ff96143c9" "checksum base64 0.9.2 (registry+https://github.com/rust-lang/crates.io-index)" = "85415d2594767338a74a30c1d370b2f3262ec1b4ed2d7bba5b3faf4de40467d9" +"checksum bitflags 0.9.1 (registry+https://github.com/rust-lang/crates.io-index)" = "4efd02e230a02e18f92fc2735f44597385ed02ad8f831e7c1c1156ee5e1ab3a5" "checksum bitflags 1.0.3 (registry+https://github.com/rust-lang/crates.io-index)" = "d0c54bb8f454c567f21197eefcdbf5679d0bd99f2ddbe52e84c77061952e6789" "checksum byteorder 1.2.3 (registry+https://github.com/rust-lang/crates.io-index)" = "74c0b906e9446b0a2e4f760cdb3fa4b2c48cdc6db8766a845c54b6ff063fd2e9" "checksum bytes 0.4.8 (registry+https://github.com/rust-lang/crates.io-index)" = "7dd32989a66957d3f0cba6588f15d4281a733f4e9ffc43fcd2385f57d3bf99ff" @@ -1288,6 +1535,8 @@ dependencies = [ "checksum cfg-if 0.1.3 (registry+https://github.com/rust-lang/crates.io-index)" = "405216fd8fe65f718daa7102ea808a946b6ce40c742998fbfd3463645552de18" "checksum chrono 0.4.3 (registry+https://github.com/rust-lang/crates.io-index)" = "a81892f0d5a53f46fc05ef0b917305a81c13f1f13bb59ac91ff595817f0764b1" "checksum clap 2.31.2 (registry+https://github.com/rust-lang/crates.io-index)" = "f0f16b89cbb9ee36d87483dc939fe9f1e13c05898d56d7b230a0d4dff033a536" +"checksum core-foundation 0.2.3 (registry+https://github.com/rust-lang/crates.io-index)" = "25bfd746d203017f7d5cbd31ee5d8e17f94b6521c7af77ece6c9e4b2d4b16c67" +"checksum core-foundation-sys 0.2.3 (registry+https://github.com/rust-lang/crates.io-index)" = "065a5d7ffdcbc8fa145d6f0746f3555025b9097a9e9cda59f7467abae670c78d" "checksum crossbeam-deque 0.3.1 (registry+https://github.com/rust-lang/crates.io-index)" = "fe8153ef04a7594ded05b427ffad46ddeaf22e63fd48d42b3e1e3bb4db07cae7" "checksum crossbeam-epoch 0.4.1 (registry+https://github.com/rust-lang/crates.io-index)" = "9b4e2817eb773f770dcb294127c011e22771899c21d18fce7dd739c0b9832e81" "checksum crossbeam-utils 0.3.2 (registry+https://github.com/rust-lang/crates.io-index)" = "d636a8b3bcc1b409d7ffd3facef8f21dcb4009626adbd0c5e6c4305c07253c7b" @@ -1296,6 +1545,8 @@ dependencies = [ "checksum errno 0.2.3 (registry+https://github.com/rust-lang/crates.io-index)" = "b2c858c42ac0b88532f48fca88b0ed947cad4f1f64d904bcd6c9f138f7b95d70" "checksum error-chain 0.11.0 (registry+https://github.com/rust-lang/crates.io-index)" = "ff511d5dc435d703f4971bc399647c9bc38e20cb41452e3b9feb4765419ed3f3" "checksum filetime 0.2.1 (registry+https://github.com/rust-lang/crates.io-index)" = "da4b9849e77b13195302c174324b5ba73eec9b236b24c221a61000daefb95c5f" +"checksum foreign-types 0.3.2 (registry+https://github.com/rust-lang/crates.io-index)" = "f6f339eb8adc052cd2ca78910fda869aefa38d22d5cb648e6485e4d3fc06f3b1" +"checksum foreign-types-shared 0.1.1 (registry+https://github.com/rust-lang/crates.io-index)" = "00b0228411908ca8685dba7fc2cdd70ec9990a6e753e89b6ac91a84c40fbaf4b" "checksum fs_extra 1.1.0 (registry+https://github.com/rust-lang/crates.io-index)" = "5f2a4a2034423744d2cc7ca2068453168dcdb82c438419e639a26bd87839c674" "checksum fuchsia-zircon 0.3.3 (registry+https://github.com/rust-lang/crates.io-index)" = "2e9763c69ebaae630ba35f74888db465e49e259ba1bc0eda7d06f4a067615d82" "checksum fuchsia-zircon-sys 0.3.3 (registry+https://github.com/rust-lang/crates.io-index)" = "3dcaa9ae7725d12cdb85b3ad99a434db70b468c09ded17e012d86b5c1010f7a7" @@ -1303,11 +1554,14 @@ dependencies = [ "checksum futures-cpupool 0.1.8 (registry+https://github.com/rust-lang/crates.io-index)" = "ab90cde24b3319636588d0c35fe03b1333857621051837ed769faefb4c2162e4" "checksum httparse 1.2.4 (registry+https://github.com/rust-lang/crates.io-index)" = "c2f407128745b78abc95c0ffbe4e5d37427fdc0d45470710cfef8c44522a2e37" "checksum humantime 1.1.1 (registry+https://github.com/rust-lang/crates.io-index)" = "0484fda3e7007f2a4a0d9c3a703ca38c71c54c55602ce4660c419fd32e188c9e" +"checksum hyper 0.10.13 (registry+https://github.com/rust-lang/crates.io-index)" = "368cb56b2740ebf4230520e2b90ebb0461e69034d85d1945febd9b3971426db2" "checksum hyper 0.11.27 (registry+https://github.com/rust-lang/crates.io-index)" = "34a590ca09d341e94cddf8e5af0bbccde205d5fbc2fa3c09dd67c7f85cea59d7" +"checksum idna 0.1.5 (registry+https://github.com/rust-lang/crates.io-index)" = "38f09e0f0b1fb55fdee1f17470ad800da77af5186a1a76c026b679358b7e844e" "checksum iovec 0.1.2 (registry+https://github.com/rust-lang/crates.io-index)" = "dbe6e417e7d0975db6512b90796e8ce223145ac4e33c377e4a42882a0e88bb08" "checksum itoa 0.4.1 (registry+https://github.com/rust-lang/crates.io-index)" = "c069bbec61e1ca5a596166e55dfe4773ff745c3d16b700013bcaff9a6df2c682" "checksum kernel32-sys 0.2.2 (registry+https://github.com/rust-lang/crates.io-index)" = "7507624b29483431c0ba2d82aece8ca6cdba9382bff4ddd0f7490560c056098d" "checksum language-tags 0.2.2 (registry+https://github.com/rust-lang/crates.io-index)" = "a91d884b6667cd606bb5a69aa0c99ba811a115fc68915e7056ec08a46e93199a" +"checksum lazy_static 0.2.11 (registry+https://github.com/rust-lang/crates.io-index)" = "76f033c7ad61445c5b347c7382dd1237847eb1bce590fe50365dcb33d546be73" "checksum lazy_static 1.0.1 (registry+https://github.com/rust-lang/crates.io-index)" = "e6412c5e2ad9584b0b8e979393122026cdd6d2a80b933f890dcd694ddbe73739" "checksum lazycell 0.6.0 (registry+https://github.com/rust-lang/crates.io-index)" = "a6f08839bc70ef4a3fe1d566d5350f519c5912ea86be0df1740a7d247c7fc0ef" "checksum libc 0.2.42 (registry+https://github.com/rust-lang/crates.io-index)" = "b685088df2b950fccadf07a7187c8ef846a959c142338a48f9dc0b94517eb5f1" @@ -1316,21 +1570,26 @@ dependencies = [ "checksum log 0.3.9 (registry+https://github.com/rust-lang/crates.io-index)" = "e19e8d5c34a3e0e2223db8e060f9e8264aeeb5c5fc64a4ee9965c062211c024b" "checksum log 0.4.2 (registry+https://github.com/rust-lang/crates.io-index)" = "6fddaa003a65722a7fb9e26b0ce95921fe4ba590542ced664d8ce2fa26f9f3ac" "checksum lru-cache 0.1.1 (registry+https://github.com/rust-lang/crates.io-index)" = "4d06ff7ff06f729ce5f4e227876cb88d10bc59cd4ae1e09fbb2bde15c850dc21" +"checksum matches 0.1.7 (registry+https://github.com/rust-lang/crates.io-index)" = "835511bab37c34c47da5cb44844bea2cfde0236db0b506f90ea4224482c9774a" "checksum memchr 2.0.1 (registry+https://github.com/rust-lang/crates.io-index)" = "796fba70e76612589ed2ce7f45282f5af869e0fdd7cc6199fa1aa1f1d591ba9d" "checksum memmap 0.6.2 (registry+https://github.com/rust-lang/crates.io-index)" = "e2ffa2c986de11a9df78620c01eeaaf27d94d3ff02bf81bfcca953102dd0c6ff" "checksum memoffset 0.2.1 (registry+https://github.com/rust-lang/crates.io-index)" = "0f9dc261e2b62d7a622bf416ea3c5245cdd5d9a7fcc428c0d06804dfce1775b3" +"checksum mime 0.2.6 (registry+https://github.com/rust-lang/crates.io-index)" = "ba626b8a6de5da682e1caa06bdb42a335aee5a84db8e5046a3e8ab17ba0a3ae0" "checksum mime 0.3.7 (registry+https://github.com/rust-lang/crates.io-index)" = "0b28683d0b09bbc20be1c9b3f6f24854efb1356ffcffee08ea3f6e65596e85fa" "checksum mio 0.6.14 (registry+https://github.com/rust-lang/crates.io-index)" = "6d771e3ef92d58a8da8df7d6976bfca9371ed1de6619d9d5a5ce5b1f29b85bfe" "checksum mio-named-pipes 0.1.6 (registry+https://github.com/rust-lang/crates.io-index)" = "f5e374eff525ce1c5b7687c4cef63943e7686524a387933ad27ca7ec43779cb3" "checksum mio-uds 0.6.6 (registry+https://github.com/rust-lang/crates.io-index)" = "84c7b5caa3a118a6e34dbac36504503b1e8dc5835e833306b9d6af0e05929f79" "checksum miow 0.2.1 (registry+https://github.com/rust-lang/crates.io-index)" = "8c1f2f3b1cf331de6896aabf6e9d55dca90356cc9960cca7eaaf408a355ae919" "checksum miow 0.3.1 (registry+https://github.com/rust-lang/crates.io-index)" = "9224c91f82b3c47cf53dcf78dfaa20d6888fbcc5d272d5f2fcdf8a697f3c987d" +"checksum native-tls 0.1.5 (registry+https://github.com/rust-lang/crates.io-index)" = "f74dbadc8b43df7864539cedb7bc91345e532fdd913cfdc23ad94f4d2d40fbc0" "checksum net2 0.2.32 (registry+https://github.com/rust-lang/crates.io-index)" = "9044faf1413a1057267be51b5afba8eb1090bd2231c693664aa1db716fe1eae0" "checksum nix 0.11.0 (registry+https://github.com/rust-lang/crates.io-index)" = "d37e713a259ff641624b6cb20e3b12b2952313ba36b6823c0f16e6cfd9e5de17" "checksum nodrop 0.1.12 (registry+https://github.com/rust-lang/crates.io-index)" = "9a2228dca57108069a5262f2ed8bd2e82496d2e074a06d1ccc7ce1687b6ae0a2" "checksum num-integer 0.1.38 (registry+https://github.com/rust-lang/crates.io-index)" = "6ac0ea58d64a89d9d6b7688031b3be9358d6c919badcf7fbb0527ccfd891ee45" "checksum num-traits 0.2.4 (registry+https://github.com/rust-lang/crates.io-index)" = "775393e285254d2f5004596d69bb8bc1149754570dcc08cf30cabeba67955e28" "checksum num_cpus 1.8.0 (registry+https://github.com/rust-lang/crates.io-index)" = "c51a3322e4bca9d212ad9a158a02abc6934d005490c054a2778df73a70aa0a30" +"checksum openssl 0.9.24 (registry+https://github.com/rust-lang/crates.io-index)" = "a3605c298474a3aa69de92d21139fb5e2a81688d308262359d85cdd0d12a7985" +"checksum openssl-sys 0.9.35 (registry+https://github.com/rust-lang/crates.io-index)" = "912f301a749394e1025d9dcddef6106ddee9252620e6d0a0e5f8d0681de9b129" "checksum percent-encoding 1.0.1 (registry+https://github.com/rust-lang/crates.io-index)" = "31010dd2e1ac33d5b46a5b413495239882813e0369f8ed8a5e266f173602f831" "checksum pkg-config 0.3.11 (registry+https://github.com/rust-lang/crates.io-index)" = "110d5ee3593dbb73f56294327fe5668bcc997897097cbc76b51e7aed3f52452f" "checksum proc-macro2 0.4.6 (registry+https://github.com/rust-lang/crates.io-index)" = "effdb53b25cdad54f8f48843d67398f7ef2e14f12c1b4cb4effc549a6462a4d6" @@ -1348,13 +1607,17 @@ dependencies = [ "checksum rustc-demangle 0.1.8 (registry+https://github.com/rust-lang/crates.io-index)" = "76d7ba1feafada44f2d38eed812bd2489a03c0f5abb975799251518b68848649" "checksum safemem 0.2.0 (registry+https://github.com/rust-lang/crates.io-index)" = "e27a8b19b835f7aea908818e871f5cc3a5a186550c30773be987e155e8163d8f" "checksum same-file 1.0.2 (registry+https://github.com/rust-lang/crates.io-index)" = "cfb6eded0b06a0b512c8ddbcf04089138c9b4362c2f696f3c3d76039d68f3637" +"checksum schannel 0.1.13 (registry+https://github.com/rust-lang/crates.io-index)" = "dc1fabf2a7b6483a141426e1afd09ad543520a77ac49bd03c286e7696ccfd77f" "checksum scoped-tls 0.1.2 (registry+https://github.com/rust-lang/crates.io-index)" = "332ffa32bf586782a3efaeb58f127980944bbc8c4d6913a86107ac2a5ab24b28" "checksum scopeguard 0.3.3 (registry+https://github.com/rust-lang/crates.io-index)" = "94258f53601af11e6a49f722422f6e3425c52b06245a5cf9bc09908b174f5e27" +"checksum security-framework 0.1.16 (registry+https://github.com/rust-lang/crates.io-index)" = "dfa44ee9c54ce5eecc9de7d5acbad112ee58755239381f687e564004ba4a2332" +"checksum security-framework-sys 0.1.16 (registry+https://github.com/rust-lang/crates.io-index)" = "5421621e836278a0b139268f36eee0dc7e389b784dc3f79d8f11aabadf41bead" "checksum serde 1.0.66 (registry+https://github.com/rust-lang/crates.io-index)" = "e9a2d9a9ac5120e0f768801ca2b58ad6eec929dc9d1d616c162f208869c2ce95" "checksum serde_bytes 0.10.4 (registry+https://github.com/rust-lang/crates.io-index)" = "adb6e51a6b3696b301bc221d785f898b4457c619b51d7ce195a6d20baecb37b3" "checksum serde_cbor 0.8.2 (registry+https://github.com/rust-lang/crates.io-index)" = "b4ad7872ff6e6c2a9221f4c1abe681e7eefc56ca5b3e87196afbfc717d141dc8" "checksum serde_derive 1.0.66 (registry+https://github.com/rust-lang/crates.io-index)" = "0a90213fa7e0f5eac3f7afe2d5ff6b088af515052cc7303bd68c7e3b91a3fb79" "checksum serde_json 1.0.20 (registry+https://github.com/rust-lang/crates.io-index)" = "fc97cccc2959f39984524026d760c08ef0dd5f0f5948c8d31797dbfae458c875" +"checksum sha1 0.2.0 (registry+https://github.com/rust-lang/crates.io-index)" = "cc30b1e1e8c40c121ca33b86c23308a090d19974ef001b4bf6e61fd1a0fb095c" "checksum slab 0.3.0 (registry+https://github.com/rust-lang/crates.io-index)" = "17b4fcaed89ab08ef143da37bc52adbcc04d4a69014f4c1208d6b51f0c47bc23" "checksum slab 0.4.0 (registry+https://github.com/rust-lang/crates.io-index)" = "fdeff4cd9ecff59ec7e3744cbca73dfe5ac35c2aedb2cfba8a1c715a18912e9d" "checksum smallvec 0.2.1 (registry+https://github.com/rust-lang/crates.io-index)" = "4c8cbcd6df1e117c2210e13ab5109635ad68a929fcbb8964dc965b76cb5ee013" @@ -1384,15 +1647,22 @@ dependencies = [ "checksum tokio-tcp 0.1.0 (registry+https://github.com/rust-lang/crates.io-index)" = "ec9b094851aadd2caf83ba3ad8e8c4ce65a42104f7b94d9e6550023f0407853f" "checksum tokio-threadpool 0.1.4 (registry+https://github.com/rust-lang/crates.io-index)" = "b3c3873a6d8d0b636e024e77b9a82eaab6739578a06189ecd0e731c7308fbc5d" "checksum tokio-timer 0.2.4 (registry+https://github.com/rust-lang/crates.io-index)" = "028b94314065b90f026a21826cffd62a4e40a92cda3e5c069cc7b02e5945f5e9" +"checksum tokio-tls 0.1.4 (registry+https://github.com/rust-lang/crates.io-index)" = "772f4b04e560117fe3b0a53e490c16ddc8ba6ec437015d91fa385564996ed913" "checksum tokio-udp 0.1.0 (registry+https://github.com/rust-lang/crates.io-index)" = "137bda266504893ac4774e0ec4c2108f7ccdbcb7ac8dced6305fe9e4e0b5041a" "checksum tokio-uds 0.1.7 (registry+https://github.com/rust-lang/crates.io-index)" = "65ae5d255ce739e8537221ed2942e0445f4b3b813daebac1c0050ddaaa3587f9" "checksum toml 0.4.6 (registry+https://github.com/rust-lang/crates.io-index)" = "a0263c6c02c4db6c8f7681f9fd35e90de799ebd4cfdeab77a38f4ff6b3d8c0d9" +"checksum traitobject 0.1.0 (registry+https://github.com/rust-lang/crates.io-index)" = "efd1f82c56340fdf16f2a953d7bda4f8fdffba13d93b00844c25572110b26079" "checksum try-lock 0.1.0 (registry+https://github.com/rust-lang/crates.io-index)" = "ee2aa4715743892880f70885373966c83d73ef1b0838a664ef0c76fffd35e7c2" +"checksum typeable 0.1.2 (registry+https://github.com/rust-lang/crates.io-index)" = "1410f6f91f21d1612654e7cc69193b0334f909dcf2c790c4826254fbb86f8887" "checksum ucd-util 0.1.1 (registry+https://github.com/rust-lang/crates.io-index)" = "fd2be2d6639d0f8fe6cdda291ad456e23629558d466e2789d2c3e9892bda285d" +"checksum unicase 1.4.2 (registry+https://github.com/rust-lang/crates.io-index)" = "7f4765f83163b74f957c797ad9253caf97f103fb064d3999aea9568d09fc8a33" "checksum unicase 2.1.0 (registry+https://github.com/rust-lang/crates.io-index)" = "284b6d3db520d67fbe88fd778c21510d1b0ba4a551e5d0fbb023d33405f6de8a" +"checksum unicode-bidi 0.3.4 (registry+https://github.com/rust-lang/crates.io-index)" = "49f2bd0c6468a8230e1db229cff8029217cf623c767ea5d60bfbd42729ea54d5" +"checksum unicode-normalization 0.1.7 (registry+https://github.com/rust-lang/crates.io-index)" = "6a0180bc61fc5a987082bfa111f4cc95c4caff7f9799f3e46df09163a937aa25" "checksum unicode-width 0.1.5 (registry+https://github.com/rust-lang/crates.io-index)" = "882386231c45df4700b275c7ff55b6f3698780a650026380e72dabe76fa46526" "checksum unicode-xid 0.1.0 (registry+https://github.com/rust-lang/crates.io-index)" = "fc72304796d0818e357ead4e000d19c9c174ab23dc11093ac919054d20a6a7fc" "checksum unreachable 1.0.0 (registry+https://github.com/rust-lang/crates.io-index)" = "382810877fe448991dfc7f0dd6e3ae5d58088fd0ea5e35189655f84e6814fa56" +"checksum url 1.7.1 (registry+https://github.com/rust-lang/crates.io-index)" = "2a321979c09843d272956e73700d12c4e7d3d92b2ee112b31548aef0d4efc5a6" "checksum utf8-ranges 1.0.0 (registry+https://github.com/rust-lang/crates.io-index)" = "662fab6525a98beff2921d7f61a39e7d59e0b425ebc7d0d9e66d316e55124122" "checksum vcpkg 0.2.3 (registry+https://github.com/rust-lang/crates.io-index)" = "7ed0f6789c8a85ca41bbc1c9d175422116a9869bd1cf31bb08e1493ecce60380" "checksum vec_map 0.8.1 (registry+https://github.com/rust-lang/crates.io-index)" = "05c78687fb1a80548ae3250346c3db86a80a7cdd77bda190189f2d0a0987c81a" @@ -1400,6 +1670,7 @@ dependencies = [ "checksum void 1.0.2 (registry+https://github.com/rust-lang/crates.io-index)" = "6a02e4885ed3bc0f2de90ea6dd45ebcbb66dacffe03547fadbb0eeae2770887d" "checksum walkdir 2.1.4 (registry+https://github.com/rust-lang/crates.io-index)" = "63636bd0eb3d00ccb8b9036381b526efac53caf112b7783b730ab3f8e44da369" "checksum want 0.0.4 (registry+https://github.com/rust-lang/crates.io-index)" = "a05d9d966753fa4b5c8db73fcab5eed4549cfe0e1e4e66911e5564a0085c35d1" +"checksum websocket 0.20.3 (registry+https://github.com/rust-lang/crates.io-index)" = "9234b4e667c19995475227172446884f516ec0963380afa960d962ab9f4c0bfa" "checksum winapi 0.2.8 (registry+https://github.com/rust-lang/crates.io-index)" = "167dc9d6949a9b857f3451275e911c3f44255842c1f7a76f33c55103a909087a" "checksum winapi 0.3.5 (registry+https://github.com/rust-lang/crates.io-index)" = "773ef9dcc5f24b7d850d0ff101e542ff24c3b090a9768e03ff889fdef41f00fd" "checksum winapi-build 0.1.1 (registry+https://github.com/rust-lang/crates.io-index)" = "2d315eee3b34aca4797b2da6b13ed88266e6d612562a0c46390af8299fc699bc" diff --git a/python/rain/capnp b/python/rain/capnp deleted file mode 120000 index 20b2c2b..0000000 --- a/python/rain/capnp +++ /dev/null @@ -1 +0,0 @@ -../../rain_core/capnp/ \ No newline at end of file diff --git a/python/rain/client/client.py b/python/rain/client/client.py index 9422b32..e9c4b30 100644 --- a/python/rain/client/client.py +++ b/python/rain/client/client.py @@ -1,24 +1,30 @@ -import capnp import json -from . import rpc -from ..common import RainException, SessionException, TaskException -from ..common.attributes import ObjectInfo, TaskInfo -from ..common.data_instance import DataInstance -from ..common.ids import governor_id_from_capnp, id_from_capnp, id_to_capnp +from .rpc import WsCommunicator, ALL_TASKS_ID from .data import DataObject from .session import Session from .task import Task +from ..common import RainException, SessionException, TaskException +from ..common.attributes import ObjectInfo, TaskInfo +from ..common.data_instance import DataInstance +from ..common.ids import ID CLIENT_PROTOCOL_VERSION = 1 FETCH_SIZE = 8 << 20 # 8MB def check_result(sessions, result): - if result.which() == "ok": + if result is None: + return + + status = result if isinstance(result, str) else result["status"] + + if status == "Ok": return # Do nothing - elif result.which() == "error": - task_id = id_from_capnp(result.error.task) + elif isinstance(status, list) and status[0] == "Error": + data = status[1] + + task_id = ID._from_json(data["task"]) message = [] if task_id.session_id == -1: @@ -41,13 +47,13 @@ def check_result(sessions, result): message.append("Task {} failed".format(task)) - message.append("Message: " + result.error.message) + message.append("Message: " + data["message"]) if task: message.append("Task created at:\n" + task._stack) - if result.error.debug: - message.append("Debug:\n" + result.error.debug) + if data["debug"]: + message.append("Debug:\n" + data["debug"]) message = "\n".join(message) raise cls(message) else: @@ -61,12 +67,10 @@ class Client: """ def __init__(self, address, port): - self._rpc_client = capnp.TwoPartyClient("{}:{}".format(address, port)) - - bootstrap = self._rpc_client.bootstrap().cast_as( - rpc.server.ServerBootstrap) - registration = bootstrap.registerAsClient(CLIENT_PROTOCOL_VERSION) - self._service = registration.wait().service + self._rpc_client = WsCommunicator(address, port) + self._rpc_client.request("RegisterClient", { + "version": CLIENT_PROTOCOL_VERSION + }) def new_session(self, name="Unnamed Session", default=False): """ @@ -78,7 +82,9 @@ def new_session(self, name="Unnamed Session", default=False): :class:`Session`: A new session """ spec = json.dumps({"name": str(name)}) - session_id = self._service.newSession(spec).wait().sessionId + session_id = self._rpc_client.request("NewSession", { + "spec": spec + })["session_id"] return Session(self, session_id, default) def get_server_info(self): @@ -88,31 +94,30 @@ def get_server_info(self): Returns: dict: A JSON-like dictionary. """ - info = self._service.getServerInfo().wait() + info = self._rpc_client.request("GetServerInfo")["governors"] return { - "governors": [{"governor_id": governor_id_from_capnp(w.governorId), - "tasks": [id_from_capnp(t) for t in w.tasks], - "objects": [id_from_capnp(o) for o in w.objects], - "objects_to_delete": [id_from_capnp(o) for o in w.objectsToDelete], - "resources": {"cpus": w.resources.nCpus}} - for w in info.governors] + "governors": [{ + "governor_id": g["governor_id"], + "tasks": [ID._from_json(id) for id in g["tasks"]], + "objects": [ID._from_json(id) for id in g["objects"]], + "objects_to_delete": [ID._from_json(id) for id in g["objects_to_delete"]], + "resources": g["resources"], + } for g in info] } def _submit(self, tasks, dataobjs): - req = self._service.submit_request() - - # Serialize tasks print(tasks, dataobjs) - - req.init("tasks", len(tasks)) - for i in range(len(tasks)): - req.tasks[i].spec = json.dumps(tasks[i].spec._to_json()) + # Serialize tasks + tasks_data = [{ + "spec": json.dumps(t.spec._to_json()) + } for t in tasks] # Serialize objects - req.init("objects", len(dataobjs)) - for i in range(len(dataobjs)): - dataobjs[i]._to_capnp(req.objects[i]) + objects_data = [d._to_json() for d in dataobjs] - req.send().wait() + return self._rpc_client.request("Submit", { + "tasks": tasks_data, + "objects": objects_data + }) def _fetch(self, dataobj): "Fetch the object data and update its state." @@ -124,30 +129,34 @@ def _fetch(self, dataobj): raise RainException( "Object {} is not submitted.".format(dataobj)) - req = self._service.fetch_request() - id_to_capnp(dataobj.id, req.id) - req.offset = 0 - req.size = FETCH_SIZE - req.includeInfo = True - result = req.send().wait() - check_result((dataobj._session,), result.status) + msg = { + "id": dataobj.id, + "include_info": True, + "offset": 0, + "size": FETCH_SIZE + } - dataobj._info = ObjectInfo._from_json(json.loads(result.info)) + result = self._rpc_client.request("Fetch", msg) + check_result((dataobj._session,), result) - size = result.transportSize - offset = len(result.data) - data = [result.data] + dataobj._info = ObjectInfo._from_json(json.loads(result["info"])) + + size = result["transport_size"] + offset = len(result["data"]) + data = [bytearray(result["data"])] while offset < size: - req = self._service.fetch_request() - id_to_capnp(dataobj.id, req.id) - req.offset = offset - req.size = FETCH_SIZE - req.includeInfo = False - r = req.send().wait() - check_result((dataobj._session,), r.status) - data.append(r.data) - offset += len(r.data) + msg = { + "id": dataobj.id, + "include_info": False, + "offset": offset, + "size": FETCH_SIZE + } + + result = self._rpc_client.request("Fetch", msg) + check_result((dataobj._session,), result["status"]) + data.append(result["data"]) + offset += len(result["data"]) rawdata = b"".join(data) return DataInstance(data=rawdata, @@ -155,67 +164,64 @@ def _fetch(self, dataobj): data_type=dataobj.spec.data_type) def _wait(self, tasks, dataobjs): - req = self._service.wait_request() - - req.init("taskIds", len(tasks)) sessions = [] for i in range(len(tasks)): task = tasks[i] if task.state is None: raise RainException("Task {} is not submitted".format(task)) - id_to_capnp(task.id, req.taskIds[i]) sessions.append(task._session) - req.init("objectIds", len(dataobjs)) + msg = { + "task_ids": [t.id for t in tasks], + "object_ids": [d.id for d in dataobjs] + } + for i in range(len(dataobjs)): - id_to_capnp(dataobjs[i].id, req.objectIds[i]) sessions.append(dataobjs[i]._session) - result = req.send().wait() + result = self._rpc_client.request("Wait", msg) check_result(sessions, result) def _close_session(self, session): - self._service.closeSession(session.session_id).wait() + return self._rpc_client.request("CloseSession", { + "session_id": session.session_id + }, allow_failure=True) def _wait_some(self, tasks, dataobjs): - req = self._service.waitSome_request() - tasks_dict = {} - req.init("taskIds", len(tasks)) for i in range(len(tasks)): tasks_dict[tasks[i].id] = tasks[i] - id_to_capnp(tasks[i].id, req.taskIds[i]) dataobjs_dict = {} - req.init("objectIds", len(dataobjs)) for i in range(len(dataobjs)): dataobjs_dict[dataobjs[i].id] = dataobjs[i] - id_to_capnp(dataobjs[i].id, req.objectIds[i]) - finished = req.send().wait() - finished_tasks = [tasks_dict[f_task.id] - for f_task in finished.finishedTasks] - finished_dataobjs = [dataobjs_dict[f_dataobj.id] - for f_dataobj in finished.finishedObjects] + msg = { + "task_ids": [t.id for t in tasks], + "object_ids": [d.id for d in dataobjs] + } + + finished = self._rpc_client.request("WaitSome", msg) + finished_tasks = [tasks_dict[ID._from_json(f_task["id"])] + for f_task in finished["finished_tasks"]] + finished_dataobjs = [dataobjs_dict[ID._from_json(f_dataobj["id"])] + for f_dataobj in finished["finished_objects"]] return finished_tasks, finished_dataobjs def _wait_all(self, session): - req = self._service.wait_request() - req.init("taskIds", 1) - req.taskIds[0].id = rpc.common.allTasksId - req.taskIds[0].sessionId = session.session_id - result = req.send().wait() + msg = { + "task_ids": [ID(session_id=session.session_id, id=ALL_TASKS_ID)], + "object_ids": [] + } + + result = self._rpc_client.request("Wait", msg) check_result((session,), result) def _unkeep(self, dataobjs): - req = self._service.unkeep_request() - - req.init("objectIds", len(dataobjs)) - for i in range(len(dataobjs)): - id_to_capnp(dataobjs[i].id, req.objectIds[i]) - - result = req.send().wait() + result = self._rpc_client.request("Unkeep", { + "object_ids": [d.id for d in dataobjs] + }, allow_failure=True) check_result([o._session for o in dataobjs], result) def update(self, items): @@ -223,31 +229,31 @@ def update(self, items): self._get_state(tasks, dataobjects) def _get_state(self, tasks, dataobjs): - req = self._service.getState_request() sessions = [] - req.init("taskIds", len(tasks)) for i in range(len(tasks)): - id_to_capnp(tasks[i].id, req.taskIds[i]) sessions.append(tasks[i]._session) dataobjs_dict = {} - req.init("objectIds", len(dataobjs)) for i in range(len(dataobjs)): dataobjs_dict[dataobjs[i].id.id] = dataobjs[i] - id_to_capnp(dataobjs[i].id, req.objectIds[i]) sessions.append(dataobjs[i]._session) - results = req.send().wait() - check_result(sessions, results.state) + msg = { + "task_ids": [t.id for t in tasks], + "object_ids": [d.id for d in dataobjs] + } + + results = self._rpc_client.request("GetState", msg)["update"] + check_result(sessions, results) - for task_update, task in zip(results.tasks, tasks): - task._state = task_update.state - task._info = TaskInfo._from_json(json.loads(task_update.info)) + for task_update, task in zip(results["tasks"], tasks): + task._state = task_update["state"] + task._info = TaskInfo._from_json(json.loads(task_update["info"])) - for object_update in results.objects: - dataobj = dataobjs_dict[object_update.id.id] - dataobj._state = object_update.state - dataobj._info = ObjectInfo._from_json(json.loads(object_update.info)) + for object_update in results["objects"]: + dataobj = dataobjs_dict[ID._from_json(object_update["id"]).id] + dataobj._state = object_update["state"] + dataobj._info = ObjectInfo._from_json(json.loads(object_update["info"])) def split_items(items): diff --git a/python/rain/client/data.py b/python/rain/client/data.py index f041936..89e9f32 100644 --- a/python/rain/client/data.py +++ b/python/rain/client/data.py @@ -2,8 +2,6 @@ import json import tarfile -import capnp - from ..common import ID, DataType, RainException from ..common.attributes import ObjectSpec from ..common.content_type import (check_content_type, encode_value, @@ -83,15 +81,13 @@ def is_kept(self): """Returns the value of self._keep""" return self._keep - def _to_capnp(self, out): - out.spec = json.dumps(self._spec._to_json()) - out.keep = self._keep - - if self._data is not None: - out.data = self._data - out.hasData = True - else: - out.hasData = False + def _to_json(self): + return { + "spec": json.dumps(self._spec._to_json()), + "keep": self._keep, + "has_data": self._data is not None, + "data": b'' if not self._data else self._data + } def wait(self): self._session.wait((self,)) @@ -115,12 +111,7 @@ def update(self): def __del__(self): if self.state is not None and self._keep: - try: - self._session.client._unkeep((self,)) - except capnp.lib.capnp.KjException: - # Ignore capnp exception, since this constructor may be - # called when connection is closed - pass + self._session.client._unkeep((self,)) def __reduce__(self): """Speciaization to replace with executor.unpickle_input_object diff --git a/python/rain/client/rpc.py b/python/rain/client/rpc.py index 3a6b6be..8812de0 100644 --- a/python/rain/client/rpc.py +++ b/python/rain/client/rpc.py @@ -1,4 +1,74 @@ -from ..common.fs import load_capnp +import asyncio -common = load_capnp("common.capnp") -server = load_capnp("server.capnp") +import cbor +import websockets + +ALL_TASKS_ID = -2 +ALL_DATA_OBJECTS_ID = -2 + + +class TaskState(object): + Finished = "Finished" + NotAssigned = "NotAssigned" + + +class DataObjectState(object): + Finished = "Finished" + Unfinished = "Unfinished" + + +def block(fut, loop=None): + if not loop: + loop = asyncio.get_event_loop() + + return loop.run_until_complete(fut) + + +class WsCommunicator(object): + def __init__(self, address, port): + self.id = 0 + self.ws = block(websockets.connect("ws://{}:{}".format(address, port), + subprotocols=['rain-ws'], + max_size=None)) + + def request(self, method, data=None, allow_failure=False): + if not self.ws.open: + if allow_failure: + return + else: + raise Exception("Client is not connected") + + if data is None: + data = {} + + id = self.id + self.id += 1 + + @asyncio.coroutine + def receive(): + try: + msg = yield from self.ws.recv() + msg = self.deserialize(msg) + + if msg["id"] == id: + return msg["data"][1] + else: + print("Error: invalid id received, sent: {}, received: {}" + .format(id, msg["id"])) + except websockets.ConnectionClosed as e: + if not allow_failure: + raise e + + msg = { + "id": id, + "data": [method, data] + } + + block(self.ws.send(self.serialize(msg))) + return block(receive()) + + def serialize(self, msg): + return cbor.dumps(msg) + + def deserialize(self, msg): + return cbor.loads(msg) diff --git a/python/rain/client/session.py b/python/rain/client/session.py index d1a0c1d..3eb6564 100644 --- a/python/rain/client/session.py +++ b/python/rain/client/session.py @@ -169,10 +169,10 @@ def submit(self): """"Submit all unsubmitted objects.""" self.client._submit(self._tasks, self._dataobjs) for task in self._tasks: - task._state = rpc.common.TaskState.notAssigned + task._state = rpc.TaskState.NotAssigned self._submitted_tasks.append(task) for dataobj in self._dataobjs: - dataobj._state = rpc.common.DataObjectState.unfinished + dataobj._state = rpc.DataObjectState.Unfinished self._submitted_dataobjs.append(dataobj) self._tasks = [] self._dataobjs = [] @@ -199,10 +199,10 @@ def wait(self, items): self.client._wait(tasks, dataobjs) for task in tasks: - task._state = rpc.common.TaskState.finished + task._state = rpc.TaskState.Finished for dataobj in dataobjs: - dataobj._state = rpc.common.DataObjectState.finished + dataobj._state = rpc.DataObjectState.Finished def wait_some(self, items): """Wait until *some* of specified tasks/dataobjects are finished. @@ -214,10 +214,10 @@ def wait_some(self, items): tasks, dataobjs) for task in finished_tasks: - task._state = rpc.common.TaskState.finished + task._state = rpc.TaskState.Finished for dataobj in finished_dataobjs: - dataobj._state = rpc.common.DataObjectState.finished + dataobj._state = rpc.DataObjectState.Finished return finished_tasks, finished_dataobjs @@ -226,10 +226,10 @@ def wait_all(self): self.client._wait_all(self) for task in self._submitted_tasks: - task._state = rpc.common.TaskState.finished + task._state = rpc.TaskState.Finished for dataobj in self._submitted_dataobjs: - dataobj._state = rpc.common.DataObjectState.finished + dataobj._state = rpc.DataObjectState.Finished def fetch(self, dataobject): """Wait for the object to finish, update its state and diff --git a/python/rain/common/data_instance.py b/python/rain/common/data_instance.py index 9bceb31..655388e 100644 --- a/python/rain/common/data_instance.py +++ b/python/rain/common/data_instance.py @@ -164,35 +164,6 @@ def write(self, path): f = tarfile.open(fileobj=io.BytesIO(self._data)) f.extractall(path) - # def _to_capnp(self, builder): - # "Internal serializer." - # if self._object_id: - # builder.storage.init("inGovernor") - # id_to_capnp(self._object_id, builder.storage.inGovernor) - # elif self._path: - # builder.storage.path = self._path - # else: - # builder.storage.memory = self._data - # attributes_to_capnp(self.attributes, builder.attributes) - - # @classmethod - # def _from_capnp(cls, reader): - # "Internal deserializer user ." - # which = reader.storage.which() - # data = None - # path = None - # if which == "memory": - # data = reader.storage.memory - # elif which == "path": - # path = reader.storage.path - # else: - # raise Exception("Invalid storage type") - # attributes = attributes_from_capnp(reader.attributes) - # return cls(data=data, - # path=path, - # attributes=attributes, - # data_type=DataType.from_capnp(reader.dataType)) - def __repr__(self): if self._data: return "".format(format_size(len(self._data)), self.attributes) diff --git a/python/rain/common/fs.py b/python/rain/common/fs.py index 376e665..e2c34d9 100644 --- a/python/rain/common/fs.py +++ b/python/rain/common/fs.py @@ -1,8 +1,6 @@ import os import shutil -import capnp - def remove_dir_content(path): """Remove content of the directory but not the directory itself""" @@ -25,9 +23,3 @@ def fresh_copy_dir(source_path, target_path): fresh_copy_dir(s, t) else: shutil.copyfile(s, t) - - -def load_capnp(filename): - src_dir = os.path.dirname(__file__) - capnp.remove_import_hook() - return capnp.load(os.path.join(src_dir, "../capnp", filename)) diff --git a/python/rain/common/ids.py b/python/rain/common/ids.py index 8eb7e2c..baf8eda 100644 --- a/python/rain/common/ids.py +++ b/python/rain/common/ids.py @@ -16,22 +16,3 @@ def _from_json(cls, data): def _to_json(self): return [self[0], self[1]] - - -def id_from_capnp(reader): - return ID(session_id=reader.sessionId, id=reader.id) - - -def id_to_capnp(obj, builder): - builder.sessionId = obj.session_id - builder.id = obj.id - - -def governor_id_from_capnp(reader): - if reader.address.which() == "ipv4": - address = reader.address.ipv4 - elif reader.address.which() == "ipv6": - raise Exception("Not implemented") - else: - raise Exception("Unknown address") - return "{}:{}".format(".".join(map(str, address)), reader.port) diff --git a/python/requirements.txt b/python/requirements.txt index 4267c2d..4998838 100644 --- a/python/requirements.txt +++ b/python/requirements.txt @@ -1,4 +1,4 @@ cbor cloudpickle pyarrow -pycapnp +websockets diff --git a/python/setup.py b/python/setup.py index 96aaeda..53142da 100644 --- a/python/setup.py +++ b/python/setup.py @@ -44,5 +44,4 @@ def load_version(): author_email='rain@substantic.net', license='MIT', packages=find_packages(), - package_data={'rain': ['capnp/*.capnp']}, install_requires=requirements) diff --git a/rain_core/Cargo.toml b/rain_core/Cargo.toml index 01cdae3..54f78a3 100644 --- a/rain_core/Cargo.toml +++ b/rain_core/Cargo.toml @@ -39,6 +39,7 @@ serde_derive = "1.0" serde_json = "1.0" tokio-core="0.1" tokio-timer = "0.2" +websocket = "0.20.3" [build-dependencies] capnpc = "0.8" diff --git a/rain_core/build.rs b/rain_core/build.rs index 3016fc4..e4caaeb 100644 --- a/rain_core/build.rs +++ b/rain_core/build.rs @@ -4,7 +4,6 @@ fn main() { capnpc::CompilerCommand::new() .file("capnp/common.capnp") .file("capnp/server.capnp") - .file("capnp/client.capnp") .file("capnp/governor.capnp") .file("capnp/monitor.capnp") .run() diff --git a/rain_core/capnp/client.capnp b/rain_core/capnp/client.capnp deleted file mode 100644 index ea32b14..0000000 --- a/rain_core/capnp/client.capnp +++ /dev/null @@ -1,99 +0,0 @@ -@0xb3195a92eff52478; - -using import "common.capnp".TaskId; -using import "common.capnp".GovernorId; -using import "common.capnp".DataObjectId; -using import "common.capnp".SessionId; -using import "common.capnp".TaskState; -using import "common.capnp".DataObjectState; -using import "common.capnp".UnitResult; -using import "common.capnp".Resources; -using import "common.capnp".DataType; -using import "common.capnp".FetchResult; - -struct GovernorInfo { - governorId @0: GovernorId; - tasks @1 :List(TaskId); - objects @2 :List(DataObjectId); - objectsToDelete @3 :List(DataObjectId); - resources @4 :Resources; -} - -struct ServerInfo { - governors @0 :List(GovernorInfo); -} - -interface ClientService { - getServerInfo @0 () -> ServerInfo; - # Get information about server - - newSession @1 (spec: Text) -> (sessionId: SessionId); - # Ask for a new session - - closeSession @2 (sessionId :SessionId) -> (); - # Remove session from governor, all running tasks are stopped, - # all existing data objects are removed - - submit @3 (tasks :List(Task), objects :List(DataObject)) -> (); - # Submit new tasks and data objects into server - # allTaskId / allDataObjectsId is NOT allowed - - unkeep @4 (objectIds :List(DataObjectId)) -> UnitResult; - # Removed "keep" flag from data objects - # It is an error if called for non-keep object - # allDataObjectsId is allowed - - wait @5 (taskIds :List(TaskId), objectIds: List(DataObjectId)) -> UnitResult; - # Wait until all given data objects are not produced - # and all given task finished. - # allTaskId / allDataObjectsId is allowed - - waitSome @6 (taskIds: List(TaskId), - objectIds: List(DataObjectId)) -> ( - finishedTasks: List(TaskId), - finishedObjects: List(DataObjectId)); - # Wait until at least one data object or task is not finished. - # It may return more objects/tasks at once. - # finished_tasks and finished_objects are both returned empty - # only if taskIds and objectsIds are empty. - # allTaskId / allDataObjectsId is allowed - - getState @7 (taskIds: List(TaskId), - objectIds: List(DataObjectId)) -> Update; - # Get current state of tasks and objects - # allTaskId / allDataObjectsId is allowed - - terminateServer @8 () -> (); - # Quit server; the connection to the server will be closed after this call - - fetch @9 (id :DataObjectId, includeInfo :Bool, offset :UInt64, size :UInt64) -> FetchResult; -} - -struct Update { - tasks @0 :List(TaskUpdate); - objects @1 :List(DataObjectUpdate); - state @2 :UnitResult; - - struct TaskUpdate { - id @0 :TaskId; - state @1 :TaskState; - info @2 :Text; - } - - struct DataObjectUpdate { - id @0 :DataObjectId; - state @1 :DataObjectState; - info @2 :Text; - } -} - -struct Task { - spec @0: Text; -} - -struct DataObject { - spec @0: Text; - keep @1 :Bool; - hasData @2: Bool; - data @3 :Data; -} diff --git a/rain_core/capnp/server.capnp b/rain_core/capnp/server.capnp index 984e8d2..31cda76 100644 --- a/rain_core/capnp/server.capnp +++ b/rain_core/capnp/server.capnp @@ -1,6 +1,5 @@ @0xb01bcb96f4bd00be; -using import "client.capnp".ClientService; using import "governor.capnp".GovernorControl; using import "governor.capnp".GovernorUpstream; using import "common.capnp".SocketAddress; @@ -8,10 +7,7 @@ using import "common.capnp".GovernorId; using import "common.capnp".Resources; interface ServerBootstrap { - registerAsClient @0 (version :Int32) -> (service :ClientService); - # Registers as a client, verifies the API version and returns the Client interface. - - registerAsGovernor @1 (version :Int32, + registerAsGovernor @0 (version :Int32, address :SocketAddress, control: GovernorControl, resources: Resources) diff --git a/rain_core/src/comm/client_message.rs b/rain_core/src/comm/client_message.rs new file mode 100644 index 0000000..2fb3f15 --- /dev/null +++ b/rain_core/src/comm/client_message.rs @@ -0,0 +1,236 @@ +use serde::Serializer; +use types::{DataObjectId, GovernorId, Resources, SessionId, TaskId}; + +#[derive(Serialize, Deserialize)] +pub struct ClientToServerMessage { + pub id: u32, + pub data: RequestType, +} +#[derive(Serialize, Deserialize)] +pub struct ServerToClientMessage { + pub id: u32, + pub data: ResponseType, +} + +#[derive(Serialize, Deserialize)] +pub enum RequestType { + RegisterClient(RegisterClientRequest), + NewSession(NewSessionRequest), + CloseSession(CloseSessionRequest), + GetServerInfo(GetServerInfoRequest), + Submit(SubmitRequest), + Fetch(FetchRequest), + Unkeep(UnkeepRequest), + Wait(WaitRequest), + WaitSome(WaitSomeRequest), + GetState(GetStateRequest), + TerminateServer(TerminateServerRequest), +} +#[derive(Serialize, Deserialize)] +pub enum ResponseType { + RegisterClient(RegisterClientResponse), + NewSession(NewSessionResponse), + CloseSession(CloseSessionResponse), + GetServerInfo(GetServerInfoResponse), + Submit(SubmitResponse), + Fetch(FetchResponse), + Unkeep(UnkeepResponse), + Wait(WaitResponse), + WaitSome(WaitSomeResponse), + GetState(GetStateResponse), + TerminateServer(TerminateServerResponse), +} + +// common types +#[derive(Serialize, Deserialize)] +pub struct RpcError { + pub message: String, + pub debug: String, + pub task: TaskId, +} +#[derive(Serialize, Deserialize)] +pub enum RpcResult { + Ok, + Error(RpcError), +} +#[derive(Serialize, Deserialize)] +pub struct Update { + pub tasks: Vec, + pub objects: Vec, + pub status: RpcResult, +} +#[derive(Serialize, Deserialize)] +pub struct TaskUpdate { + pub id: TaskId, + pub state: TaskState, + pub info: String, +} +#[derive(Serialize, Deserialize)] +pub struct DataObjectUpdate { + pub id: DataObjectId, + pub state: DataObjectState, + pub info: String, +} +#[derive(Serialize, Deserialize)] +pub enum TaskState { + NotAssigned, + Ready, + Assigned, + Running, + Finished, + Failed, +} +#[derive(Serialize, Deserialize)] +pub enum DataObjectState { + Unfinished, + Finished, + Removed, +} + +// request/response types +#[derive(Serialize, Deserialize)] +pub struct RegisterClientRequest { + pub version: u32, +} +#[derive(Serialize, Deserialize)] +pub struct RegisterClientResponse {} + +#[derive(Serialize, Deserialize)] +pub struct NewSessionRequest { + pub spec: String, +} +#[derive(Serialize, Deserialize)] +pub struct NewSessionResponse { + pub session_id: SessionId, +} + +#[derive(Serialize, Deserialize)] +pub struct CloseSessionRequest { + pub session_id: SessionId, +} +#[derive(Serialize, Deserialize)] +pub struct CloseSessionResponse {} + +#[derive(Serialize, Deserialize)] +pub struct GetServerInfoRequest {} +#[derive(Serialize, Deserialize)] +pub struct GetServerInfoResponse { + pub governors: Vec, +} +#[derive(Serialize, Deserialize)] +pub struct GovernorInfo { + #[serde(serialize_with = "serialize_socket_addr")] + pub governor_id: GovernorId, + pub tasks: Vec, + pub objects: Vec, + pub objects_to_delete: Vec, + pub resources: Resources, +} + +fn serialize_socket_addr(addr: &GovernorId, s: S) -> Result +where + S: Serializer, +{ + s.serialize_str(&format!("{}", addr)) +} + +#[derive(Serialize, Deserialize)] +pub struct SubmitRequest { + pub tasks: Vec, + pub objects: Vec, +} +#[derive(Serialize, Deserialize)] +pub struct SubmitResponse {} +#[derive(Serialize, Deserialize)] +pub struct Task { + pub spec: String, +} +#[derive(Serialize, Deserialize)] +pub struct DataObject { + pub spec: String, + pub keep: bool, + #[serde(with = "::serde_bytes")] + pub data: Vec, + pub has_data: bool, +} + +#[derive(Serialize, Deserialize)] +pub struct FetchRequest { + pub id: DataObjectId, + pub include_info: bool, + pub offset: u64, + pub size: u64, +} +#[derive(Serialize, Deserialize)] +pub struct FetchResponse { + pub status: FetchStatus, + #[serde(with = "::serde_bytes")] + pub data: Vec, + pub info: String, + pub transport_size: u64, +} + +impl FetchResponse { + pub fn error(error: RpcError) -> Self { + FetchResponse { + status: FetchStatus::Error(error), + data: vec![], + info: "".to_owned(), + transport_size: 0, + } + } +} +#[derive(Serialize, Deserialize)] +pub enum FetchStatus { + Ok, + Redirect(GovernorId), + NotHere, + Removed, + Error(RpcError), + Ignored, +} + +#[derive(Serialize, Deserialize)] +pub struct UnkeepRequest { + pub object_ids: Vec, +} +#[derive(Serialize, Deserialize)] +pub struct UnkeepResponse { + pub status: RpcResult, +} + +#[derive(Serialize, Deserialize)] +pub struct WaitRequest { + pub task_ids: Vec, + pub object_ids: Vec, +} +#[derive(Serialize, Deserialize)] +pub struct WaitResponse { + pub status: RpcResult, +} + +#[derive(Serialize, Deserialize)] +pub struct WaitSomeRequest { + pub task_ids: Vec, + pub object_ids: Vec, +} +#[derive(Serialize, Deserialize)] +pub struct WaitSomeResponse { + pub finished_tasks: Vec, + pub finished_objects: Vec, +} + +#[derive(Serialize, Deserialize)] +pub struct GetStateRequest { + pub task_ids: Vec, + pub object_ids: Vec, +} +#[derive(Serialize, Deserialize)] +pub struct GetStateResponse { + pub update: Update, +} + +#[derive(Serialize, Deserialize)] +pub struct TerminateServerRequest {} +#[derive(Serialize, Deserialize)] +pub struct TerminateServerResponse {} diff --git a/rain_core/src/comm/mod.rs b/rain_core/src/comm/mod.rs index cac5179..92a3c90 100644 --- a/rain_core/src/comm/mod.rs +++ b/rain_core/src/comm/mod.rs @@ -1,5 +1,7 @@ +pub mod client_message; pub(crate) mod executor; -pub use self::executor::{CallMsg, DataLocation, DropCachedMsg, ExecutorToGovernorMessage, - GovernorToExecutorMessage, LocalObjectIn, LocalObjectOut, RegisterMsg, - ResultMsg}; +pub use self::executor::{ + CallMsg, DataLocation, DropCachedMsg, ExecutorToGovernorMessage, GovernorToExecutorMessage, + LocalObjectIn, LocalObjectOut, RegisterMsg, ResultMsg, +}; diff --git a/rain_core/src/errors.rs b/rain_core/src/errors.rs index fd817d7..a542b84 100644 --- a/rain_core/src/errors.rs +++ b/rain_core/src/errors.rs @@ -18,7 +18,9 @@ error_chain!{ SessionErr(SessionError); Utf8Err(::std::str::Utf8Error); Json(::serde_json::Error); + Cbor(::serde_cbor::error::Error); Sqlite(::rusqlite::Error); + Websocket(::websocket::WebSocketError); } errors { @@ -39,9 +41,9 @@ impl ::std::convert::From for ::capnp::Error { #[derive(Debug, Clone)] pub struct SessionError { - message: String, - debug: String, - task_id: TaskId, + pub message: String, + pub debug: String, + pub task_id: TaskId, } impl SessionError { diff --git a/rain_core/src/lib.rs b/rain_core/src/lib.rs index 0f9a165..68cbe5c 100644 --- a/rain_core/src/lib.rs +++ b/rain_core/src/lib.rs @@ -31,6 +31,7 @@ extern crate serde_derive; extern crate serde_json; extern crate tokio_core; extern crate tokio_timer; +extern crate websocket; pub const VERSION: &str = env!("CARGO_PKG_VERSION"); pub const GOVERNOR_PROTOCOL_VERSION: i32 = 0; @@ -50,10 +51,6 @@ pub mod server_capnp { include!(concat!(env!("OUT_DIR"), "/capnp/server_capnp.rs")); } -pub mod client_capnp { - include!(concat!(env!("OUT_DIR"), "/capnp/client_capnp.rs")); -} - pub mod common_capnp { use std::fmt; diff --git a/rain_core/src/types/id.rs b/rain_core/src/types/id.rs index 8a036d1..a2223d9 100644 --- a/rain_core/src/types/id.rs +++ b/rain_core/src/types/id.rs @@ -81,17 +81,14 @@ pub trait SId: for<'a> ToCapnp<'a> + for<'a> FromCapnp<'a> + WriteCapnp + ReadCa /// ID type for task objects. #[derive(Copy, Clone, Debug, Ord, Eq, PartialEq, PartialOrd, Hash, Default)] pub struct TaskId { - session_id: SessionId, - id: Id, + pub session_id: SessionId, + pub id: Id, } impl SId for TaskId { #[inline] fn new(session_id: SessionId, id: Id) -> Self { - TaskId { - session_id: session_id, - id: id, - } + TaskId { session_id, id } } #[inline] @@ -146,17 +143,14 @@ impl<'a> FromCapnp<'a> for TaskId { /// ID type for task objects. #[derive(Copy, Clone, Debug, Ord, Eq, PartialEq, PartialOrd, Hash, Default)] pub struct DataObjectId { - session_id: SessionId, - id: Id, + pub session_id: SessionId, + pub id: Id, } impl SId for DataObjectId { #[inline] fn new(session_id: SessionId, id: Id) -> Self { - DataObjectId { - session_id: session_id, - id: id, - } + DataObjectId { session_id, id } } #[inline] diff --git a/rain_server/Cargo.toml b/rain_server/Cargo.toml index 73b1442..bac6360 100644 --- a/rain_server/Cargo.toml +++ b/rain_server/Cargo.toml @@ -56,6 +56,7 @@ tokio-timer = "0.2" tokio-uds="0.1" toml = "0.4" walkdir = "2" +websocket = "0.20.3" [build-dependencies] capnpc = "0.8" diff --git a/rain_server/src/main.rs b/rain_server/src/main.rs index 3904dd6..f4d3c37 100644 --- a/rain_server/src/main.rs +++ b/rain_server/src/main.rs @@ -39,6 +39,7 @@ extern crate tokio_timer; extern crate tokio_uds; extern crate toml; extern crate walkdir; +extern crate websocket; extern crate rain_core; @@ -62,7 +63,8 @@ use rain_core::sys::{create_ready_file, get_hostname}; use rain_core::{errors::*, utils::*}; pub const VERSION: &str = env!("CARGO_PKG_VERSION"); -const DEFAULT_SERVER_PORT: u16 = 7210; +const DEFAULT_SERVER_GOVERNOR_PORT: u16 = 7210; +const DEFAULT_SERVER_CLIENT_PORT: u16 = 7211; const DEFAULT_GOVERNOR_PORT: u16 = 0; const DEFAULT_HTTP_SERVER_PORT: u16 = 8080; @@ -81,13 +83,23 @@ fn parse_listen_arg(key: &str, args: &ArgMatches, default_port: u16) -> SocketAd } fn run_server(_global_args: &ArgMatches, cmd_args: &ArgMatches) { - let listen_address = parse_listen_arg("LISTEN_ADDRESS", cmd_args, DEFAULT_SERVER_PORT); + let governor_listen_address = parse_listen_arg( + "GOVERNOR_LISTEN_ADDRESS", + cmd_args, + DEFAULT_SERVER_GOVERNOR_PORT, + ); + let client_listen_address = parse_listen_arg( + "CLIENT_LISTEN_ADDRESS", + cmd_args, + DEFAULT_SERVER_CLIENT_PORT, + ); let http_listen_address = parse_listen_arg("HTTP_LISTEN_ADDRESS", cmd_args, DEFAULT_HTTP_SERVER_PORT); let ready_file = cmd_args.value_of("READY_FILE"); info!("Starting Rain {} server", VERSION); - info!("Listen address: {}", listen_address); + info!("Governor listen address: {}", governor_listen_address); + info!("Client listen address: {}", client_listen_address); let log_dir = cmd_args .value_of("LOG_DIR") @@ -120,7 +132,8 @@ fn run_server(_global_args: &ArgMatches, cmd_args: &ArgMatches) { let state = server::state::StateRef::new( tokio_core.handle(), - listen_address, + governor_listen_address, + client_listen_address, http_listen_address, log_dir, test_mode, @@ -199,7 +212,7 @@ fn run_governor(_global_args: &ArgMatches, cmd_args: &ArgMatches) { let mut server_address = cmd_args.value_of("SERVER_ADDRESS").unwrap().to_string(); if !server_address.contains(':') { - server_address = format!("{}:{}", server_address, DEFAULT_SERVER_PORT); + server_address = format!("{}:{}", server_address, DEFAULT_SERVER_GOVERNOR_PORT); } let server_addr = match server_address.to_socket_addrs() { @@ -324,7 +337,16 @@ fn run_governor(_global_args: &ArgMatches, cmd_args: &ArgMatches) { } fn run_starter(_global_args: &ArgMatches, cmd_args: &ArgMatches) { - let listen_address = parse_listen_arg("LISTEN_ADDRESS", cmd_args, DEFAULT_SERVER_PORT); + let governor_listen_address = parse_listen_arg( + "GOVERNOR_LISTEN_ADDRESS", + cmd_args, + DEFAULT_SERVER_GOVERNOR_PORT, + ); + let client_listen_address = parse_listen_arg( + "CLIENT_LISTEN_ADDRESS", + cmd_args, + DEFAULT_SERVER_CLIENT_PORT, + ); let http_listen_address = parse_listen_arg("HTTP_LISTEN_ADDRESS", cmd_args, DEFAULT_HTTP_SERVER_PORT); let log_dir = cmd_args @@ -375,7 +397,8 @@ fn run_starter(_global_args: &ArgMatches, cmd_args: &ArgMatches) { let mut config = start::starter::StarterConfig::new( local_governors, - listen_address, + governor_listen_address, + client_listen_address, http_listen_address, &log_dir, cmd_args.value_of("REMOTE_INIT").unwrap_or("").to_string(), @@ -467,10 +490,15 @@ fn main() { .subcommand( // ---- SERVER ---- SubCommand::with_name("server") .about("Rain server") - .arg(Arg::with_name("LISTEN_ADDRESS") - .short("l") - .long("--listen") - .help("Listening port/address/address:port (default 0.0.0.0:7210)") + .arg(Arg::with_name("GOVERNOR_LISTEN_ADDRESS") + .short("g") + .long("--governor-listen") + .help("Governor listening port/address/address:port (default 0.0.0.0:7210)") + .takes_value(true)) + .arg(Arg::with_name("CLIENT_LISTEN_ADDRESS") + .short("c") + .long("--client-listen") + .help("Client listening port/address/address:port (default 0.0.0.0:7211)") .takes_value(true)) .arg(Arg::with_name("HTTP_LISTEN_ADDRESS") .long("--http-listen") @@ -549,11 +577,15 @@ fn main() { .arg(Arg::with_name("RCOS") // RCOS = Reserve CPUs on Server .short("-S") .help("Reserve a CPU on server machine")) - .arg(Arg::with_name("LISTEN_ADDRESS") - .short("l") - .value_name("ADDRESS") - .long("--listen") - .help("Server listening port/address/address:port (default = 0.0.0.0:auto)") + .arg(Arg::with_name("GOVERNOR_LISTEN_ADDRESS") + .short("g") + .long("--governor-listen") + .help("Governor listening port/address/address:port (default 0.0.0.0:7210)") + .takes_value(true)) + .arg(Arg::with_name("CLIENT_LISTEN_ADDRESS") + .short("c") + .long("--client-listen") + .help("Client listening port/address/address:port (default 0.0.0.0:7211)") .takes_value(true)) .arg(Arg::with_name("HTTP_LISTEN_ADDRESS") .long("--http-listen") diff --git a/rain_server/src/server/mod.rs b/rain_server/src/server/mod.rs index b1817b1..7782b0d 100644 --- a/rain_server/src/server/mod.rs +++ b/rain_server/src/server/mod.rs @@ -5,3 +5,4 @@ pub mod rpc; pub mod scheduler; pub mod state; pub mod testmode; +pub mod ws; diff --git a/rain_server/src/server/rpc/bootstrap.rs b/rain_server/src/server/rpc/bootstrap.rs index b4443c5..6c14835 100644 --- a/rain_server/src/server/rpc/bootstrap.rs +++ b/rain_server/src/server/rpc/bootstrap.rs @@ -5,10 +5,10 @@ use rain_core::server_capnp::server_bootstrap; use rain_core::{types::*, utils::*}; use std::net::SocketAddr; -use super::{ClientServiceImpl, GovernorUpstreamImpl}; +use super::GovernorUpstreamImpl; use server::state::StateRef; -use rain_core::{CLIENT_PROTOCOL_VERSION, GOVERNOR_PROTOCOL_VERSION}; +use rain_core::GOVERNOR_PROTOCOL_VERSION; // ServerBootstrap is the entry point of RPC service. // It is created on the server and provided @@ -25,7 +25,7 @@ impl ServerBootstrapImpl { ServerBootstrapImpl { state: state.clone(), registered: false, - address: address, + address, } } } @@ -37,36 +37,6 @@ impl Drop for ServerBootstrapImpl { } impl server_bootstrap::Server for ServerBootstrapImpl { - fn register_as_client( - &mut self, - params: server_bootstrap::RegisterAsClientParams, - mut results: server_bootstrap::RegisterAsClientResults, - ) -> Promise<(), ::capnp::Error> { - if self.registered { - error!("Multiple registration from connection {}", self.address); - return Promise::err(capnp::Error::failed(format!( - "Connection already registered" - ))); - } - - let params = pry!(params.get()); - - if params.get_version() != CLIENT_PROTOCOL_VERSION { - error!("Client protocol mismatch, expected {}, got {}", CLIENT_PROTOCOL_VERSION, params.get_version()); - return Promise::err(capnp::Error::failed(format!("Client protocol mismatch, expected {}, got {}", CLIENT_PROTOCOL_VERSION, params.get_version()))); - } - - self.registered = true; - - let service = ::rain_core::client_capnp::client_service::ToClient::new(pry!( - ClientServiceImpl::new(&self.state, &self.address) - )).from_server::<::capnp_rpc::Server>(); - - info!("Connection {} registered as client", self.address); - results.get().set_service(service); - Promise::ok(()) - } - fn register_as_governor( &mut self, params: server_bootstrap::RegisterAsGovernorParams, diff --git a/rain_server/src/server/rpc/client.rs b/rain_server/src/server/rpc/client.rs deleted file mode 100644 index 7405852..0000000 --- a/rain_server/src/server/rpc/client.rs +++ /dev/null @@ -1,494 +0,0 @@ -use capnp::capability::Promise; -use futures::{future, Future}; -use rain_core::client_capnp::client_service; -use rain_core::{errors::*, types::*, utils::*}; -use std::net::SocketAddr; - -use server::graph::{ClientRef, TaskRef}; -use server::graph::{DataObjectRef, DataObjectState}; -use server::state::StateRef; - -pub struct ClientServiceImpl { - state: StateRef, - client: ClientRef, -} - -impl ClientServiceImpl { - pub fn new(state: &StateRef, address: &SocketAddr) -> Result { - Ok(Self { - state: state.clone(), - client: state.get_mut().add_client(address.clone())?, - }) - } -} - -impl Drop for ClientServiceImpl { - fn drop(&mut self) { - let mut s = self.state.get_mut(); - info!("Client {} disconnected", self.client.get_id()); - s.remove_client(&self.client) - .expect("client connection drop"); - } -} - -impl client_service::Server for ClientServiceImpl { - fn get_server_info( - &mut self, - _: client_service::GetServerInfoParams, - mut results: client_service::GetServerInfoResults, - ) -> Promise<(), ::capnp::Error> { - debug!("Client asked for info"); - let s = self.state.get(); - - let futures: Vec<_> = s.graph - .governors - .iter() - .map(|(governor_id, governor)| { - let w = governor.get(); - let control = w.control.as_ref().unwrap(); - let governor_id = governor_id.clone(); - let resources = w.resources.clone(); - control - .get_info_request() - .send() - .promise - .map(move |r| (governor_id, r, resources)) - }) - .collect(); - - Promise::from_future(future::join_all(futures).map(move |rs| { - let results = results.get(); - let mut governors = results.init_governors(rs.len() as u32); - for (i, &(ref governor_id, ref r, ref resources)) in rs.iter().enumerate() { - let mut w = governors.reborrow().get(i as u32); - let r = r.get().unwrap(); - w.set_tasks(r.get_tasks().unwrap()).unwrap(); - w.set_objects(r.get_objects().unwrap()).unwrap(); - w.set_objects_to_delete(r.get_objects_to_delete().unwrap()) - .unwrap(); - resources.to_capnp(&mut w.reborrow().get_resources().unwrap()); - governor_id.to_capnp(&mut w.get_governor_id().unwrap()); - } - () - })) - } - - fn new_session( - &mut self, - params: client_service::NewSessionParams, - mut results: client_service::NewSessionResults, - ) -> Promise<(), ::capnp::Error> { - let params = pry!(params.get()); - let mut s = self.state.get_mut(); - let spec = ::serde_json::from_str(pry!(params.get_spec())).unwrap(); - let session = pry!(s.add_session(&self.client, spec)); - results.get().set_session_id(session.get_id()); - debug!("Client asked for a new session, got {:?}", session.get_id()); - Promise::ok(()) - } - - fn close_session( - &mut self, - params: client_service::CloseSessionParams, - _: client_service::CloseSessionResults, - ) -> Promise<(), ::capnp::Error> { - let params = pry!(params.get()); - let mut s = self.state.get_mut(); - let session = pry!(s.session_by_id(params.get_session_id())); - s.remove_session(&session).unwrap(); - Promise::ok(()) - } - - fn submit( - &mut self, - params: client_service::SubmitParams, - _: client_service::SubmitResults, - ) -> Promise<(), ::capnp::Error> { - let mut s = self.state.get_mut(); - let params = pry!(params.get()); - let tasks = pry!(params.get_tasks()); - let objects = pry!(params.get_objects()); - info!( - "New task submission ({} tasks, {} data objects) from client {}", - tasks.len(), - objects.len(), - self.client.get_id() - ); - debug!("Sessions: {:?}", s.graph.sessions); - let mut created_tasks = Vec::::new(); - let mut created_objects = Vec::::new(); - // catch any insertion error and clean up later - let res: Result<()> = (|| { - // first create the objects - for co in objects.iter() { - let spec: ObjectSpec = ::serde_json::from_str(co.get_spec().unwrap()).unwrap(); - let session = s.session_by_id(spec.id.get_session_id())?; - let data = if co.get_has_data() { - Some(co.get_data()?.into()) - } else { - None - }; - let o = s.add_object(&session, spec, co.get_keep(), data)?; - created_objects.push(o); - } - // second create the tasks - for ct in tasks.iter() { - let spec: TaskSpec = ::serde_json::from_str(ct.get_spec().unwrap()).unwrap(); - let session = s.session_by_id(spec.id.get_session_id())?; - let mut inputs = Vec::::with_capacity(spec.inputs.len()); - for ci in spec.inputs.iter() { - inputs.push(s.object_by_id(ci.id)?); - } - let mut outputs = Vec::::with_capacity(spec.outputs.len()); - for co in spec.outputs.iter() { - outputs.push(s.object_by_id(*co)?); - } - let t = s.add_task(&session, spec, inputs, outputs)?; - created_tasks.push(t); - } - debug!("New tasks: {:?}", created_tasks); - debug!("New objects: {:?}", created_objects); - s.logger.add_client_submit_event( - created_tasks.iter().map(|t| t.get().spec.clone()).collect(), - created_objects - .iter() - .map(|o| o.get().spec.clone()) - .collect(), - ); - // verify submit integrity - s.verify_submit(&created_tasks, &created_objects) - })(); - if res.is_err() { - debug!("Error: {:?}", res); - for t in created_tasks { - pry!(s.remove_task(&t)); - } - for o in created_objects { - pry!(s.remove_object(&o)); - } - pry!(res); - } - Promise::ok(()) - } - - fn wait( - &mut self, - params: client_service::WaitParams, - mut result: client_service::WaitResults, - ) -> Promise<(), ::capnp::Error> { - // Set error from session to result - fn set_error( - result: &mut ::rain_core::common_capnp::unit_result::Builder, - error: &SessionError, - ) { - error.to_capnp(&mut result.reborrow().init_error()); - } - - let s = self.state.get_mut(); - let params = pry!(params.get()); - let task_ids = pry!(params.get_task_ids()); - let object_ids = pry!(params.get_object_ids()); - info!( - "New wait request ({} tasks, {} data objects) from client", - task_ids.len(), - object_ids.len() - ); - - if task_ids.len() == 1 && object_ids.len() == 0 - && task_ids.get(0).get_id() == ::rain_core::common_capnp::ALL_TASKS_ID - { - let session_id = task_ids.get(0).get_session_id(); - debug!("Waiting for all session session_id={}", session_id); - let session = match s.session_by_id(session_id) { - Ok(s) => s, - Err(e) => return Promise::err(::capnp::Error::failed(e.description().to_string())), - }; - if let &Some(ref e) = session.get().get_error() { - set_error(&mut result.get(), e); - return Promise::ok(()); - } - - let session2 = session.clone(); - return Promise::from_future(session.get_mut().wait().then(move |r| { - match r { - Ok(_) => result.get().set_ok(()), - Err(_) => { - set_error( - &mut result.get(), - session2.get().get_error().as_ref().unwrap(), - ); - } - }; - Ok(()) - })); - } - - let mut sessions = RcSet::new(); - - // TODO: Wait for data objects - // TODO: Implement waiting for session (for special "all" IDs) - // TODO: Get rid of unwrap and do proper error handling - - let mut task_futures = Vec::new(); - - for id in task_ids.iter() { - match s.task_by_id_check_session(TaskId::from_capnp(&id)) { - Ok(t) => { - let mut task = t.get_mut(); - sessions.insert(task.session.clone()); - if task.is_finished() { - continue; - } - task_futures.push(task.wait()); - } - Err(Error(ErrorKind::SessionErr(ref e), _)) => { - set_error(&mut result.get(), e); - return Promise::ok(()); - } - Err(e) => return Promise::err(::capnp::Error::failed(e.description().to_string())), - }; - } - - debug!("{} waiting futures", task_futures.len()); - - if task_futures.is_empty() { - result.get().set_ok(()); - return Promise::ok(()); - } - - Promise::from_future(::futures::future::join_all(task_futures).then(move |r| { - match r { - Ok(_) => result.get().set_ok(()), - Err(_) => { - let session = sessions.iter().find(|s| s.get().is_failed()).unwrap(); - set_error( - &mut result.get(), - session.get().get_error().as_ref().unwrap(), - ); - } - }; - Ok(()) - })) - } - - fn wait_some( - &mut self, - params: client_service::WaitSomeParams, - _results: client_service::WaitSomeResults, - ) -> Promise<(), ::capnp::Error> { - let params = pry!(params.get()); - let task_ids = pry!(params.get_task_ids()); - let object_ids = pry!(params.get_object_ids()); - info!( - "New wait_some request ({} tasks, {} data objects) from client", - task_ids.len(), - object_ids.len() - ); - Promise::err(::capnp::Error::failed( - "wait_sone is not implemented yet".to_string(), - )) - } - - fn unkeep( - &mut self, - params: client_service::UnkeepParams, - mut results: client_service::UnkeepResults, - ) -> Promise<(), ::capnp::Error> { - let mut s = self.state.get_mut(); - let params = pry!(params.get()); - let object_ids = pry!(params.get_object_ids()); - debug!( - "New unkeep request ({} data objects) from client", - object_ids.len() - ); - - let mut objects = Vec::new(); - for oid in object_ids.iter() { - let id: DataObjectId = DataObjectId::from_capnp(&oid); - match s.object_by_id_check_session(id) { - Ok(obj) => objects.push(obj), - Err(Error(ErrorKind::SessionErr(ref e), _)) => { - e.to_capnp(&mut results.get().init_error()); - return Promise::ok(()); - } - Err(e) => return Promise::err(::capnp::Error::failed(e.description().to_string())), - }; - } - - for o in objects.iter() { - s.unkeep_object(&o); - } - s.logger - .add_client_unkeep_event(objects.iter().map(|o| o.get().spec.id).collect()); - Promise::ok(()) - } - - fn fetch( - &mut self, - params: client_service::FetchParams, - mut results: client_service::FetchResults, - ) -> Promise<(), ::capnp::Error> { - let params = pry!(params.get()); - let id = DataObjectId::from_capnp(&pry!(params.get_id())); - - debug!("Client fetch for object id={}", id); - - let object = match self.state.get().object_by_id_check_session(id) { - Ok(t) => t, - Err(Error(ErrorKind::SessionErr(ref e), _)) => { - e.to_capnp(&mut results.get().get_status().init_error()); - return Promise::ok(()); - } - Err(e) => return Promise::err(::capnp::Error::failed(e.description().to_string())), - }; - let object2 = object.clone(); - let mut obj = object2.get_mut(); - if obj.state == DataObjectState::Removed { - return Promise::err(::capnp::Error::failed(format!( - "create_reader on removed object {:?}", - obj - ))); - } - - let size = params.get_size(); - - if size > 32 << 20 - /* 32 MB */ - { - let mut err = results.get().get_status().init_error(); - err.set_message("Fetch size is too big."); - return Promise::ok(()); - } - - let offset = params.get_offset(); - let include_info = params.get_include_info(); - let session = obj.session.clone(); - let state_ref = self.state.clone(); - - Promise::from_future( - obj.wait() - .then(move |r| -> future::Either<_, _> { - if r.is_err() { - let session = session.get(); - session - .get_error() - .as_ref() - .unwrap() - .to_capnp(&mut results.get().get_status().init_error()); - return future::Either::A(future::result(Ok(()))); - } - let obj = object.get(); - if obj.state == DataObjectState::Removed { - let session = session.get(); - session - .get_error() - .as_ref() - .unwrap() - .to_capnp(&mut results.get().get_status().init_error()); - return future::Either::A(future::result(Ok(()))); - } - assert_eq!( - obj.state, - DataObjectState::Finished, - "triggered finish hook on unfinished object" - ); - - if obj.data.is_some() { - // Fetching uploaded objects is not implemented yet - unimplemented!(); - } - let governor_ref = obj.located.iter().next().unwrap().clone(); - let mut governor = governor_ref.get_mut(); - debug!("Redirecting client fetch id={} to {}", governor.id(), id); - future::Either::B( - governor - .wait_for_data_connection(&governor_ref, &state_ref) - .and_then(move |data_conn| { - let mut req = data_conn.fetch_request(); - { - let mut request = req.get(); - request.set_offset(offset); - request.set_size(size); - request.set_include_info(include_info); - id.to_capnp(&mut request.get_id().unwrap()); - } - req.send() - .promise - .map(move |r| { - results.set(r.get().unwrap()).unwrap(); - }) - .map_err(|e| e.into()) - }), - ) - }) - .map_err(|e| panic!("Fetch failed: {:?}", e)), - ) - } - - fn get_state( - &mut self, - params: client_service::GetStateParams, - mut results: client_service::GetStateResults, - ) -> Promise<(), ::capnp::Error> { - let params = pry!(params.get()); - let task_ids = pry!(params.get_task_ids()); - let object_ids = pry!(params.get_object_ids()); - info!( - "New get_state request ({} tasks, {} data objects) from client", - task_ids.len(), - object_ids.len() - ); - - let s = self.state.get(); - let tasks: Vec<_> = match task_ids - .iter() - .map(|id| s.task_by_id_check_session(TaskId::from_capnp(&id))) - .collect() - { - Ok(tasks) => tasks, - Err(Error(ErrorKind::SessionErr(ref e), _)) => { - e.to_capnp(&mut results.get().get_state().unwrap().init_error()); - return Promise::ok(()); - } - Err(e) => return Promise::err(::capnp::Error::failed(e.description().to_string())), - }; - - let objects: Vec<_> = match object_ids - .iter() - .map(|id| s.object_by_id_check_session(DataObjectId::from_capnp(&id))) - .collect() - { - Ok(tasks) => tasks, - Err(Error(ErrorKind::SessionErr(ref e), _)) => { - e.to_capnp(&mut results.get().get_state().unwrap().init_error()); - return Promise::ok(()); - } - Err(e) => return Promise::err(::capnp::Error::failed(e.description().to_string())), - }; - - let mut results = results.get(); - - { - let mut task_updates = results.reborrow().init_tasks(tasks.len() as u32); - for (i, task) in tasks.iter().enumerate() { - let mut update = task_updates.reborrow().get(i as u32); - let t = task.get(); - update.set_info(&::serde_json::to_string(&t.info).unwrap()); - t.spec.id.to_capnp(&mut update.get_id().unwrap()); - } - } - - { - let mut obj_updates = results.reborrow().init_objects(objects.len() as u32); - for (i, obj) in objects.iter().enumerate() { - let mut update = obj_updates.reborrow().get(i as u32); - let o = obj.get(); - update.set_info(&::serde_json::to_string(&o.info).unwrap()); - o.spec.id.to_capnp(&mut update.get_id().unwrap()); - } - } - - results.get_state().unwrap().set_ok(()); - Promise::ok(()) - } -} diff --git a/rain_server/src/server/rpc/mod.rs b/rain_server/src/server/rpc/mod.rs index dd08b9b..215f865 100644 --- a/rain_server/src/server/rpc/mod.rs +++ b/rain_server/src/server/rpc/mod.rs @@ -1,7 +1,5 @@ mod bootstrap; -mod client; mod governor; pub use self::bootstrap::ServerBootstrapImpl; -pub use self::client::ClientServiceImpl; pub use self::governor::GovernorUpstreamImpl; diff --git a/rain_server/src/server/state.rs b/rain_server/src/server/state.rs index 7e21912..0c34040 100644 --- a/rain_server/src/server/state.rs +++ b/rain_server/src/server/state.rs @@ -3,22 +3,24 @@ use std::net::SocketAddr; use std::path::PathBuf; use std::time::{Duration, Instant}; -use rain_core::logging::events; use futures::{Future, Stream}; use hyper::server::Http; +use rain_core::logging::events; use rain_core::{errors::*, sys::*, types::*, utils::*}; use tokio_core::net::{TcpListener, TcpStream}; use tokio_core::reactor::Handle; use common::new_rpc_system; -use server::graph::{ClientRef, DataObjectRef, DataObjectState, GovernorRef, Graph, SessionRef, - TaskRef, TaskState}; +use server::graph::{ + ClientRef, DataObjectRef, DataObjectState, GovernorRef, Graph, SessionRef, TaskRef, TaskState, +}; use server::http::RequestHandler; use server::logging::logger::Logger; use server::logging::sqlite_logger::SQLiteLogger; use server::rpc::ServerBootstrapImpl; use server::scheduler::{ReactiveScheduler, UpdatedIn}; use server::testmode; +use server::ws::comm::ClientCommunicator; use wrapped::WrappedRcRefCell; const LOGGING_INTERVAL: u64 = 1; // Logging interval in seconds @@ -53,8 +55,11 @@ pub struct State { pub logger: Box, - /// Listening port and address. - listen_address: SocketAddr, + /// Governor listening port and address. + governor_listen_address: SocketAddr, + + /// Client interface + client_comm: ClientCommunicator, /// Listening port for HTTP interface http_listen_address: SocketAddr, @@ -198,7 +203,8 @@ impl State { self.logger.add_closed_session_event( session.get_id(), events::SessionClosedReason::ClientClose, - String::new()); + String::new(), + ); } // remove from graph self.graph.sessions.remove(&session.get_id()).unwrap(); @@ -231,7 +237,8 @@ impl State { self.logger.add_closed_session_event( session.get_id(), events::SessionClosedReason::Error, - cause); + cause, + ); Ok(()) } @@ -436,7 +443,8 @@ impl State { let mut co = &mut new_objects.reborrow().get(0); let o = object.get(); o.to_governor_capnp(&mut co); - let placement = o.located + let placement = o + .located .iter() .next() .map(|w| w.get().id().clone()) @@ -480,7 +488,8 @@ impl State { wref.check_consistency_opt().unwrap(); // non-recoverable // Create request - let mut req = wref.get() + let mut req = wref + .get() .control .as_ref() .unwrap() @@ -542,7 +551,8 @@ impl State { let o = input.get_mut(); if !o.assigned.contains(&wref) { // Just take first placement - let placement = o.located + let placement = o + .located .iter() .next() .map(|w| w.get().id().clone()) @@ -637,7 +647,8 @@ impl State { wref.get_mut().assigned_tasks.remove(task); self.update_task_assignment(task); - for oref in task.get() + for oref in task + .get() .outputs .iter() .map(|x| x.clone()) @@ -833,8 +844,7 @@ impl State { assert_eq!(t.state, TaskState::Assigned); t.state = state; t.info = info.clone(); - self.logger - .add_task_started_event(t.id(), info); + self.logger.add_task_started_event(t.id(), info); } TaskState::Failed => { debug!( @@ -858,10 +868,7 @@ impl State { let task_id = tref.get().spec.id; self.fail_session(&session, error_message.clone(), debug_message, task_id) .unwrap(); - self.logger.add_task_finished_event( - tref.get().id(), - info - ); + self.logger.add_task_finished_event(tref.get().id(), info); } _ => panic!( "Invalid governor {:?} task {:?} state update to {:?}", @@ -952,7 +959,8 @@ impl State { && !wref.get().scheduled_ready_tasks.is_empty() { // TODO: Prioritize older members of w.scheduled_ready_tasks (order-preserving set) - let tref = wref.get() + let tref = wref + .get() .scheduled_ready_tasks .iter() .next() @@ -1023,7 +1031,8 @@ pub type StateRef = WrappedRcRefCell; impl StateRef { pub fn new( handle: Handle, - listen_address: SocketAddr, + governor_listen_address: SocketAddr, + client_listen_address: SocketAddr, http_listen_address: SocketAddr, log_dir: PathBuf, test_mode: bool, @@ -1034,10 +1043,11 @@ impl StateRef { let s = Self::wrap(State { graph, - test_mode: test_mode, - listen_address: listen_address, - http_listen_address: http_listen_address, - handle: handle, + test_mode, + governor_listen_address, + client_comm: ClientCommunicator::new(client_listen_address), + http_listen_address, + handle, scheduler: Default::default(), underload_governors: Default::default(), updates: Default::default(), @@ -1051,10 +1061,10 @@ impl StateRef { } pub fn start(&self) { - let listen_address = self.get().listen_address; + let governor_listen_address = self.get().governor_listen_address; let http_listen_address = self.get().http_listen_address; let handle = self.get().handle.clone(); - let listener = TcpListener::bind(&listen_address, &handle).unwrap(); + let listener = TcpListener::bind(&governor_listen_address, &handle).unwrap(); let state = self.clone(); let future = listener @@ -1068,6 +1078,11 @@ impl StateRef { }); handle.spawn(future); + // ---- Start Client websocket server ---- + self.get() + .client_comm + .start(self.get().handle.clone(), self.clone()); + // ---- Start HTTP server ---- //let listener = TcpListener::bind(&http_listen_address, &handle).unwrap(); let handle1 = self.get().handle.clone(); diff --git a/rain_server/src/server/ws/client.rs b/rain_server/src/server/ws/client.rs new file mode 100644 index 0000000..d75081f --- /dev/null +++ b/rain_server/src/server/ws/client.rs @@ -0,0 +1,613 @@ +use futures::future::{self, err, join_all, ok, Future}; +use rain_core::{ + comm::client_message::{ + CloseSessionRequest, CloseSessionResponse, DataObjectState, DataObjectUpdate, FetchRequest, + FetchResponse, FetchStatus, GetServerInfoRequest, GetServerInfoResponse, GetStateRequest, + GetStateResponse, GovernorInfo, NewSessionRequest, NewSessionResponse, + RegisterClientRequest, RegisterClientResponse, RpcError, RpcResult, SubmitRequest, + SubmitResponse, TaskState, TaskUpdate, TerminateServerRequest, TerminateServerResponse, + UnkeepRequest, UnkeepResponse, Update, WaitRequest, WaitResponse, WaitSomeRequest, + WaitSomeResponse, + }, + common_capnp::DataObjectState as CapnpDataObjectState, + errors::SessionError, + types::{DataObjectId, ObjectSpec, SId, TaskId, TaskSpec}, + utils::{FromCapnp, RcSet, ToCapnp}, + Error, ErrorKind, CLIENT_PROTOCOL_VERSION, +}; +use server::{ + graph::{ClientRef, DataObjectRef, TaskRef}, + state::StateRef, +}; +use std::net::SocketAddr; + +type Response = Box>; + +pub trait ClientService { + fn register_client(&mut self, msg: RegisterClientRequest) -> Response; + fn new_session(&mut self, msg: NewSessionRequest) -> Response; + fn close_session(&mut self, msg: CloseSessionRequest) -> Response; + fn get_server_info(&mut self, msg: GetServerInfoRequest) -> Response; + fn submit(&mut self, msg: SubmitRequest) -> Response; + fn fetch(&mut self, msg: FetchRequest) -> Response; + fn unkeep(&mut self, msg: UnkeepRequest) -> Response; + fn wait(&mut self, msg: WaitRequest) -> Response; + fn wait_some(&mut self, msg: WaitSomeRequest) -> Response; + fn get_state(&mut self, msg: GetStateRequest) -> Response; + fn terminate_server( + &mut self, + msg: TerminateServerRequest, + ) -> Response; +} + +pub struct ClientServiceImpl { + address: SocketAddr, + state: StateRef, + client: ClientRef, + registered: bool, +} + +impl ClientServiceImpl { + pub fn new(address: SocketAddr, state: StateRef) -> Result { + let client = state.get_mut().add_client(address.clone())?; + Ok(ClientServiceImpl { + address, + client, + state, + registered: false, + }) + } + + fn check_registration(&self) -> Result<(), Error> { + if !self.registered { + bail!("Client not registered") + } else { + Ok(()) + } + } +} + +macro_rules! fry { + ($result:expr) => { + match $result { + Ok(res) => res, + Err(e) => return Box::new(err(e.into())), + } + }; +} + +macro_rules! from_capnp_list { + ($builder:expr, $items:ident, $obj:ident) => {{ + $builder + .$items()? + .iter() + .map(|item| $obj::from_capnp(&item)) + .collect() + }}; +} + +fn convert_task_state(state: &::rain_core::common_capnp::TaskState) -> TaskState { + match state { + ::rain_core::common_capnp::TaskState::NotAssigned => TaskState::NotAssigned, + ::rain_core::common_capnp::TaskState::Ready => TaskState::Ready, + ::rain_core::common_capnp::TaskState::Assigned => TaskState::Assigned, + ::rain_core::common_capnp::TaskState::Running => TaskState::Running, + ::rain_core::common_capnp::TaskState::Finished => TaskState::Finished, + ::rain_core::common_capnp::TaskState::Failed => TaskState::Failed, + } +} +fn convert_object_state(state: &::rain_core::common_capnp::DataObjectState) -> DataObjectState { + match state { + ::rain_core::common_capnp::DataObjectState::Unfinished => DataObjectState::Unfinished, + ::rain_core::common_capnp::DataObjectState::Finished => DataObjectState::Finished, + ::rain_core::common_capnp::DataObjectState::Removed => DataObjectState::Removed, + } +} + +fn response(response: T) -> Response { + Box::new(ok(response)) +} +fn error(error: Error) -> Response { + Box::new(err(error)) +} + +impl Drop for ClientServiceImpl { + fn drop(&mut self) { + let mut s = self.state.get_mut(); + info!("Client {} disconnected", self.client.get_id()); + s.remove_client(&self.client) + .expect("client connection drop"); + } +} + +impl ClientService for ClientServiceImpl { + fn register_client(&mut self, msg: RegisterClientRequest) -> Response { + if self.registered { + error!("Multiple registration from connection {}", self.address); + return error("Connection already registered".into()); + } + + let version = msg.version as i32; + if version != CLIENT_PROTOCOL_VERSION { + error!( + "Client protocol mismatch, expected {}, got {}", + CLIENT_PROTOCOL_VERSION, version + ); + return error( + format!( + "Client protocol mismatch, expected {}, got {}", + CLIENT_PROTOCOL_VERSION, version + ).into(), + ); + } + + info!("Connection {} registered as client", self.address); + + self.registered = true; + response(RegisterClientResponse {}) + } + fn new_session(&mut self, msg: NewSessionRequest) -> Response { + fry!(self.check_registration()); + + let mut s = self.state.get_mut(); + let spec = ::serde_json::from_str(&msg.spec).unwrap(); + let session = fry!(s.add_session(&self.client, spec)); + + debug!("Client asked for a new session, got {:?}", session.get_id()); + + response(NewSessionResponse { + session_id: session.get_id(), + }) + } + fn close_session(&mut self, msg: CloseSessionRequest) -> Response { + fry!(self.check_registration()); + + let mut s = self.state.get_mut(); + let session = fry!(s.session_by_id(msg.session_id)); + s.remove_session(&session).unwrap(); + response(CloseSessionResponse {}) + } + fn get_server_info(&mut self, _: GetServerInfoRequest) -> Response { + fry!(self.check_registration()); + + debug!("Client asked for info"); + let s = self.state.get(); + + let futures: Vec<_> = s + .graph + .governors + .iter() + .map(|(governor_id, governor)| { + let w = governor.get(); + let control = w.control.as_ref().unwrap(); + let governor_id = governor_id.clone(); + let resources = w.resources.clone(); + control + .get_info_request() + .send() + .promise + .map(move |r| (governor_id, r, resources)) + }).collect(); + + Box::new( + join_all(futures) + .and_then(move |rs| { + let mut governors = vec![]; + for &(ref governor_id, ref r, ref resources) in rs.iter() { + let r = r.get()?; + + governors.push(GovernorInfo { + governor_id: *governor_id, + tasks: from_capnp_list!(r, get_tasks, TaskId), + objects: from_capnp_list!(r, get_objects, DataObjectId), + objects_to_delete: from_capnp_list!( + r, + get_objects_to_delete, + DataObjectId + ), + resources: resources.clone(), + }); + } + Ok(GetServerInfoResponse { governors }) + }).map_err(|e| e.into()), + ) + } + fn submit(&mut self, msg: SubmitRequest) -> Response { + fry!(self.check_registration()); + + let mut s = self.state.get_mut(); + let tasks = msg.tasks; + let mut objects = msg.objects; + info!( + "New task submission ({} tasks, {} data objects) from client {}", + tasks.len(), + objects.len(), + self.client.get_id() + ); + debug!("Sessions: {:?}", s.graph.sessions); + let mut created_tasks = Vec::::new(); + let mut created_objects = Vec::::new(); + // catch any insertion error and clean up later + let res: Result<(), Error> = (|| { + // first create the objects + for co in objects.iter_mut() { + let spec: ObjectSpec = ::serde_json::from_str(&co.spec).unwrap(); + let session = s.session_by_id(spec.id.session_id)?; + + let data = if co.has_data { + Some(::std::mem::replace(&mut co.data, vec![])) + } else { + None + }; + let o = s.add_object(&session, spec, co.keep, data)?; + created_objects.push(o); + } + // second create the tasks + for ct in tasks.iter() { + let spec: TaskSpec = ::serde_json::from_str(&ct.spec).unwrap(); + let session = s.session_by_id(spec.id.get_session_id())?; + let mut inputs = Vec::::with_capacity(spec.inputs.len()); + for ci in spec.inputs.iter() { + inputs.push(s.object_by_id(ci.id)?); + } + let mut outputs = Vec::::with_capacity(spec.outputs.len()); + for co in spec.outputs.iter() { + outputs.push(s.object_by_id(*co)?); + } + let t = s.add_task(&session, spec, inputs, outputs)?; + created_tasks.push(t); + } + debug!("New tasks: {:?}", created_tasks); + debug!("New objects: {:?}", created_objects); + s.logger.add_client_submit_event( + created_tasks.iter().map(|t| t.get().spec.clone()).collect(), + created_objects + .iter() + .map(|o| o.get().spec.clone()) + .collect(), + ); + // verify submit integrity + s.verify_submit(&created_tasks, &created_objects) + })(); + if res.is_err() { + debug!("Error: {:?}", res); + for t in created_tasks { + fry!(s.remove_task(&t)); + } + for o in created_objects { + fry!(s.remove_object(&o)); + } + fry!(res); + } + response(SubmitResponse {}) + } + fn fetch(&mut self, msg: FetchRequest) -> Response { + fry!(self.check_registration()); + + let id = msg.id; + debug!("Client fetch for object id={}", id); + + let object = match self.state.get().object_by_id_check_session(id) { + Ok(t) => t, + Err(Error(ErrorKind::SessionErr(ref e), _)) => { + return response(FetchResponse::error(RpcError { + message: e.message.clone(), + debug: e.debug.clone(), + task: e.task_id.clone(), + })); + } + Err(e) => return error(e.description().into()), + }; + let object2 = object.clone(); + let mut obj = object2.get_mut(); + if obj.state == CapnpDataObjectState::Removed { + return error(format!("create_reader on removed object {:?}", obj).into()); + } + + let size = msg.size; + if size > 32 << 20 + /* 32 MB */ + { + return response(FetchResponse::error(RpcError { + message: "Fetch size is too big.".to_owned(), + debug: "".to_owned(), + task: TaskId { + id: 0, + session_id: 0, + }, + })); + } + + let offset = msg.offset; + let include_info = msg.include_info; + let session = obj.session.clone(); + let state_ref = self.state.clone(); + + Box::new( + obj.wait() + .then(move |r| -> future::Either<_, _> { + if r.is_err() { + let session = session.get(); + let error = session.get_error().as_ref().unwrap(); + return future::Either::A(future::result(Ok(FetchResponse::error( + RpcError { + message: error.message.clone(), + debug: error.debug.clone(), + task: error.task_id.clone(), + }, + )))); + } + let obj = object.get(); + if obj.state == CapnpDataObjectState::Removed { + let session = session.get(); + let error = session.get_error().as_ref().unwrap(); + return future::Either::A(future::result(Ok(FetchResponse::error( + RpcError { + message: error.message.clone(), + debug: error.debug.clone(), + task: error.task_id.clone(), + }, + )))); + } + assert_eq!( + obj.state, + CapnpDataObjectState::Finished, + "triggered finish hook on unfinished object" + ); + + if obj.data.is_some() { + // Fetching uploaded objects is not implemented yet + unimplemented!(); + } + let governor_ref = obj.located.iter().next().unwrap().clone(); + let mut governor = governor_ref.get_mut(); + debug!("Redirecting client fetch id={} to {}", governor.id(), id); + future::Either::B( + governor + .wait_for_data_connection(&governor_ref, &state_ref) + .and_then(move |data_conn| { + let mut req = data_conn.fetch_request(); + { + let mut request = req.get(); + request.set_offset(offset); + request.set_size(size); + request.set_include_info(include_info); + id.to_capnp(&mut request.get_id().unwrap()); + } + req.send() + .promise + .map(move |r| { + let result = r.get().unwrap(); + + FetchResponse { + status: FetchStatus::Ok, + data: result.get_data().unwrap().to_vec(), + info: result.get_info().unwrap().to_string(), + transport_size: result.get_transport_size(), + } + }).map_err(|e| e.into()) + }), + ) + }).map_err(|e| panic!("Fetch failed: {:?}", e)), + ) + } + fn unkeep(&mut self, msg: UnkeepRequest) -> Response { + fry!(self.check_registration()); + + let mut s = self.state.get_mut(); + let object_ids = msg.object_ids; + debug!( + "New unkeep request ({} data objects) from client", + object_ids.len() + ); + + let mut objects = Vec::new(); + for oid in object_ids.iter() { + match s.object_by_id_check_session(*oid) { + Ok(obj) => objects.push(obj), + Err(Error(ErrorKind::SessionErr(ref e), _)) => { + return response(UnkeepResponse { + status: RpcResult::Error(RpcError { + message: e.message.clone(), + debug: e.debug.clone(), + task: e.task_id.clone(), + }), + }); + } + Err(e) => return error(e.description().into()), + }; + } + + for o in objects.iter() { + s.unkeep_object(&o); + } + s.logger + .add_client_unkeep_event(objects.iter().map(|o| o.get().spec.id).collect()); + response(UnkeepResponse { + status: RpcResult::Ok, + }) + } + fn wait(&mut self, msg: WaitRequest) -> Response { + fry!(self.check_registration()); + + fn session_error(error: &SessionError) -> WaitResponse { + WaitResponse { + status: RpcResult::Error(RpcError { + message: error.message.clone(), + debug: error.debug.clone(), + task: error.task_id.clone(), + }), + } + }; + fn response_ok() -> WaitResponse { + WaitResponse { + status: RpcResult::Ok, + } + } + + let s = self.state.get_mut(); + let task_ids = msg.task_ids; + let object_ids = msg.object_ids; + info!( + "New wait request ({} tasks, {} data objects) from client", + task_ids.len(), + object_ids.len() + ); + + if task_ids.len() == 1 + && object_ids.len() == 0 + && task_ids[0].id == ::rain_core::common_capnp::ALL_TASKS_ID + { + let session_id = task_ids[0].session_id; + debug!("Waiting for all session session_id={}", session_id); + let session = match s.session_by_id(session_id) { + Ok(s) => s, + Err(e) => return error(e.description().into()), + }; + if let &Some(ref e) = session.get().get_error() { + return response(session_error(e)); + } + + let session2 = session.clone(); + return Box::new(session.get_mut().wait().then(move |r| { + ok(match r { + Ok(_) => response_ok(), + Err(_) => session_error(&session2.get().get_error().clone().unwrap()), + }) + })); + } + + let mut sessions = RcSet::new(); + + // TODO: Wait for data objects + // TODO: Implement waiting for session (for special "all" IDs) + // TODO: Get rid of unwrap and do proper error handling + + let mut task_futures = Vec::new(); + + for id in task_ids.iter() { + match s.task_by_id_check_session(*id) { + Ok(t) => { + let mut task = t.get_mut(); + sessions.insert(task.session.clone()); + if task.is_finished() { + continue; + } + task_futures.push(task.wait()); + } + Err(Error(ErrorKind::SessionErr(ref e), _)) => { + return response(session_error(e)); + } + Err(e) => return error(e.description().into()), + }; + } + + debug!("{} waiting futures", task_futures.len()); + + if task_futures.is_empty() { + return Box::new(ok(response_ok())); + } + + Box::new(join_all(task_futures).then(move |r| { + ok(match r { + Ok(_) => response_ok(), + Err(_) => { + let session = sessions.iter().find(|s| s.get().is_failed()).unwrap(); + session_error(&session.get().get_error().clone().unwrap()) + } + }) + })) + } + fn wait_some(&mut self, msg: WaitSomeRequest) -> Response { + fry!(self.check_registration()); + + let task_ids = msg.task_ids; + let object_ids = msg.object_ids; + info!( + "New wait_some request ({} tasks, {} data objects) from client", + task_ids.len(), + object_ids.len() + ); + error("wait_some is not implemented yet".into()) + } + fn get_state(&mut self, msg: GetStateRequest) -> Response { + fry!(self.check_registration()); + + fn session_error(error: &SessionError) -> GetStateResponse { + GetStateResponse { + update: Update { + status: RpcResult::Error(RpcError { + message: error.message.clone(), + debug: error.debug.clone(), + task: error.task_id.clone(), + }), + tasks: vec![], + objects: vec![], + }, + } + }; + + let task_ids = msg.task_ids; + let object_ids = msg.object_ids; + info!( + "New get_state request ({} tasks, {} data objects) from client", + task_ids.len(), + object_ids.len() + ); + + let s = self.state.get(); + let tasks: Vec<_> = match task_ids + .iter() + .map(|id| s.task_by_id_check_session(*id)) + .collect() + { + Ok(tasks) => tasks, + Err(Error(ErrorKind::SessionErr(ref e), _)) => { + return response(session_error(e)); + } + Err(e) => return error(e.description().into()), + }; + + let objects: Vec<_> = match object_ids + .iter() + .map(|id| s.object_by_id_check_session(*id)) + .collect() + { + Ok(tasks) => tasks, + Err(Error(ErrorKind::SessionErr(ref e), _)) => { + return response(session_error(e)); + } + Err(e) => return error(e.description().into()), + }; + + let mut update = Update { + tasks: vec![], + objects: vec![], + status: RpcResult::Ok, + }; + + for task in tasks.iter() { + let t = task.get(); + update.tasks.push(TaskUpdate { + id: t.id(), + state: convert_task_state(&t.state), + info: ::serde_json::to_string(&t.info).unwrap(), + }); + } + for obj in objects.iter() { + let o = obj.get(); + update.objects.push(DataObjectUpdate { + id: o.id(), + state: convert_object_state(&o.state), + info: ::serde_json::to_string(&o.info).unwrap(), + }); + } + + response(GetStateResponse { update }) + } + fn terminate_server(&mut self, _: TerminateServerRequest) -> Response { + fry!(self.check_registration()); + + response(TerminateServerResponse {}) + } +} diff --git a/rain_server/src/server/ws/comm.rs b/rain_server/src/server/ws/comm.rs new file mode 100644 index 0000000..6c4a302 --- /dev/null +++ b/rain_server/src/server/ws/comm.rs @@ -0,0 +1,135 @@ +use super::client::ClientServiceImpl; +use futures::{future::err, Future, Stream}; +use rain_core::{ + comm::client_message::{ + ClientToServerMessage, RequestType, ResponseType, ServerToClientMessage, + }, + Error, +}; +use serde_cbor::{de::from_slice, ser::to_vec}; +use server::{state::StateRef, ws::client::ClientService}; +use std::{fmt::Debug, net::SocketAddr}; +use tokio_core::reactor::Handle; +use websocket::{async::Server, server::InvalidConnection, OwnedMessage}; + +macro_rules! rpc_methods { + ( $message:expr, $client:expr, $( ($method:ident, $implementation:ident) ),* ) => { + match $message.data { + $( + RequestType::$method(data) => { + let id = $message.id; + debug!("Message from client: {:?}", stringify!($method)); + return Some(Box::new($client.$implementation(data).map(move |r| { + ServerToClientMessage { + id, + data: ResponseType::$method(r) + } + }))); + } + )* + } + } +} + +fn handle_message( + m: OwnedMessage, + client: &mut dyn ClientService, +) -> Option>> { + if let OwnedMessage::Binary(data) = m { + let message = match from_slice::(&data) { + Ok(msg) => msg, + Err(e) => { + debug!("Received invalid message from client: {:?}", e); + return Some(Box::new(err(e.into()))); + } + }; + + rpc_methods!( + message, + client, + (RegisterClient, register_client), + (NewSession, new_session), + (CloseSession, close_session), + (GetServerInfo, get_server_info), + (Submit, submit), + (Fetch, fetch), + (Unkeep, unkeep), + (Wait, wait), + (WaitSome, wait_some), + (GetState, get_state), + (TerminateServer, terminate_server) + ) + } + + None +} + +pub struct ClientCommunicator { + address: SocketAddr, +} + +impl ClientCommunicator { + pub fn new(address: SocketAddr) -> Self { + ClientCommunicator { address } + } + + pub fn start(&self, handle: Handle, state: StateRef) { + let server = Server::bind(&self.address, &handle).unwrap(); + let protocol = "rain-ws"; + + let start_handle = handle.clone(); + let handler = server + .incoming() + .map_err(|InvalidConnection { error, .. }| error) + .for_each(move |(upgrade, addr)| { + if !upgrade.protocols().iter().any(|s| s == protocol) { + spawn_future(upgrade.reject(), &handle); + return Ok(()); + } + + let mut client_impl = match ClientServiceImpl::new(addr, state.clone()) { + Ok(client) => client, + Err(_) => { + spawn_future(upgrade.reject(), &handle); + return Ok(()); + } + }; + + let future = upgrade + .use_protocol(protocol) + .accept() + .map_err(|e| e.into()) + .and_then(move |(s, _)| { + let (sink, stream) = s.split(); + stream + .take_while(|m| Ok(!m.is_close())) + .map_err(|e| e.into()) + .and_then(move |m| { + if let Some(fut) = handle_message(m, &mut client_impl) { + Some(fut.map(|res| OwnedMessage::Binary(to_vec(&res).unwrap()))) + } else { + None + } + }).filter_map(|x| x) + .forward(sink) + }); + + spawn_future(future, &handle); + Ok(()) + }); + start_handle.spawn(handler.map_err(|e| panic!("RPC error: {:?}", e))); + } +} + +fn spawn_future(f: F, handle: &Handle) +where + F: Future + 'static, + E: Debug, +{ + // errors from individual clients are logged elsewhere and silently discarded here + handle.spawn( + f.map_err(|e| { + println!("RPC error: {:?}", e); + }).map(|_| {}), + ); +} diff --git a/rain_server/src/server/ws/mod.rs b/rain_server/src/server/ws/mod.rs new file mode 100644 index 0000000..e6c9051 --- /dev/null +++ b/rain_server/src/server/ws/mod.rs @@ -0,0 +1,2 @@ +pub mod client; +pub mod comm; diff --git a/rain_server/src/start/starter.rs b/rain_server/src/start/starter.rs index 75fe49a..49ad08f 100644 --- a/rain_server/src/start/starter.rs +++ b/rain_server/src/start/starter.rs @@ -16,8 +16,11 @@ pub struct StarterConfig { /// Number of local governor that will be spawned pub local_governors: Vec>, - /// Listening address of server - pub server_listen_address: SocketAddr, + /// Listening governor address of server + pub server_governor_address: SocketAddr, + + /// Listening client address of server + pub server_client_address: SocketAddr, /// Listening address of server for HTTP connections pub server_http_listen_address: SocketAddr, @@ -41,7 +44,8 @@ pub struct StarterConfig { impl StarterConfig { pub fn new( local_governors: Vec>, - server_listen_address: SocketAddr, + server_governor_address: SocketAddr, + server_client_address: SocketAddr, server_http_listen_address: SocketAddr, log_dir: &Path, remote_init: String, @@ -50,7 +54,8 @@ impl StarterConfig { ) -> Self { Self { local_governors, - server_listen_address, + server_governor_address, + server_client_address, server_http_listen_address, log_dir: ::std::env::current_dir().unwrap().join(log_dir), // Make it absolute governor_host_file: None, @@ -187,11 +192,15 @@ impl Starter { fn start_server(&mut self) -> Result<()> { let ready_file = self.create_tmp_filename("server-ready"); let (program, program_args) = self.local_rain_command(); - let server_address = format!("{}", self.config.server_listen_address); + let server_governor_address = format!("{}", self.config.server_governor_address); + let server_client_address = format!("{}", self.config.server_client_address); let server_http_address = format!("{}", self.config.server_http_listen_address); let http_port = self.config.server_http_listen_address.port(); - info!("Starting local server ({})", server_address); + info!( + "Starting local server (governor: {}, client: {})", + server_governor_address, server_client_address + ); let log_dir = self.config.log_dir.join("server"); self.server_pid = { let process = self.spawn_process( @@ -202,8 +211,10 @@ impl Starter { .arg("server") .arg("--logdir") .arg(&log_dir) - .arg("--listen") - .arg(&server_address) + .arg("--governor-listen") + .arg(&server_governor_address) + .arg("--client-listen") + .arg(&server_client_address) .arg("--http-listen") .arg(&server_http_address) .arg("--ready-file") @@ -282,7 +293,11 @@ impl Starter { } else { get_hostname() }; - format!("{}:{}", hostname, self.config.server_listen_address.port()) + format!( + "{}:{}", + hostname, + self.config.server_governor_address.port() + ) } fn start_local_governors(&mut self) -> Result<()> { @@ -292,7 +307,8 @@ impl Starter { ); let server_address = self.server_address(true); let (program, program_args) = self.local_rain_command(); - let governors: Vec<_> = self.config + let governors: Vec<_> = self + .config .local_governors .iter() .cloned() diff --git a/tests/pytests/conftest.py b/tests/pytests/conftest.py index e3e3543..a344932 100644 --- a/tests/pytests/conftest.py +++ b/tests/pytests/conftest.py @@ -53,7 +53,8 @@ def kill_all(self): class TestEnv(Env): default_listen_port = "17010" - default_http_port = "17011" + default_client_port = "17011" + default_http_port = "17012" running_port = None def __init__(self): @@ -78,6 +79,7 @@ def start(self, n_cpus=1, listen_addr=None, listen_port=None, + client_port=None, http_port=None, governor_defs=None, delete_list_timeout=None, @@ -124,7 +126,12 @@ def start(self, else: addr = self.default_listen_port port = self.default_listen_port - self.running_port = port + + if not client_port: + client_port = self.default_client_port + client_addr = "127.0.0.1:{}".format(client_port) + + self.running_port = client_port if not http_port: http_port = self.default_http_port @@ -142,7 +149,8 @@ def start(self, args = (RAIN_BIN, "server", "--ready-file", server_ready_file, "--logdir", os.path.join(WORK_DIR, "server"), - "--listen", str(addr), + "--governor-listen", str(addr), + "--client-listen", str(client_addr), "--http-listen", str(http_port)) self.server = self.start_process("server", args, env=env) assert self.server is not None diff --git a/tests/pytests/test_client.py b/tests/pytests/test_client.py index 47aca46..4fd8a7a 100644 --- a/tests/pytests/test_client.py +++ b/tests/pytests/test_client.py @@ -125,8 +125,8 @@ def test_submit(test_env): s.submit() assert s.task_count == 0 assert s.dataobj_count == 0 - assert t1.state == rpc.common.TaskState.notAssigned - assert t2.state == rpc.common.TaskState.notAssigned + assert t1.state == rpc.TaskState.NotAssigned + assert t2.state == rpc.TaskState.NotAssigned @pytest.mark.xfail(reason="wait_some not implemented") @@ -139,8 +139,8 @@ def test_wait_some(test_env): t2 = tasks.Sleep(t1, 0.4) s.submit() finished = s.wait_some((t1,), ()) - assert t1.state == rpc.common.TaskState.finished - assert t2.state == rpc.common.TaskState.notAssigned + assert t1.state == rpc.TaskState.Finished + assert t2.state == rpc.TaskState.NotAssigned assert len(finished) == 2 assert len(finished[0]) == 1 assert len(finished[1]) == 0 @@ -157,8 +157,8 @@ def test_wait_all(test_env): t2 = tasks.Sleep(t1, 0.5) s.submit() test_env.assert_duration(0.35, 0.65, lambda: s.wait_all()) - assert t1.state == rpc.common.TaskState.finished - assert t2.state == rpc.common.TaskState.finished + assert t1.state == rpc.TaskState.Finished + assert t2.state == rpc.TaskState.Finished test_env.assert_max_duration(0.1, lambda: t2.wait()) @@ -297,9 +297,9 @@ def test_task_wait(test_env): t1 = tasks.Concat((blob("a"), blob("b"))) assert t1.state is None s.submit() - assert t1.state == rpc.common.TaskState.notAssigned + assert t1.state == rpc.TaskState.NotAssigned t1.wait() - assert t1.state == rpc.common.TaskState.finished + assert t1.state == rpc.TaskState.Finished def test_fetch_removed_object_fails(test_env): @@ -353,9 +353,9 @@ def test_dataobj_wait(test_env): o1 = t1.output assert t1.state is None s.submit() - assert o1.state == rpc.common.DataObjectState.unfinished + assert o1.state == rpc.DataObjectState.Unfinished o1.wait() - assert o1.state == rpc.common.DataObjectState.finished + assert o1.state == rpc.DataObjectState.Finished def test_fetch_outputs(test_env):