szymski

Untitled

Jan 23rd, 2019
204
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
text 18.55 KB | None | 0 0
  1. using System;
  2. using System.Collections.Generic;
  3. using System.Diagnostics;
  4. using System.IO;
  5. using System.Linq;
  6. using System.Net;
  7. using System.Net.NetworkInformation;
  8. using System.Net.Sockets;
  9. using System.Text;
  10. using System.Threading;
  11. using System.Threading.Tasks;
  12. using System.Windows;
  13.  
  14. namespace ClusterXWorker
  15. {
  16. public class Peer
  17. {
  18. public IPEndPoint EndPoint { get; set; }
  19. public string ComputerName { get; set; }
  20. public DateTime LastAlive { get; set; }
  21. public bool RunningProcess { get; set; }
  22. }
  23.  
  24. public class Message
  25. {
  26. public IPEndPoint EndPoint { get; set; }
  27. public byte[] Data { get; set; }
  28. }
  29.  
  30. public class Cluster
  31. {
  32. public int Port { get; set; } = 21370;
  33. public int Timeout { get; set; } = 10000;
  34.  
  35. private UdpClient _client;
  36.  
  37. private Task _listenerTask;
  38. private Task _senderTask;
  39. private Task _discoverTask;
  40.  
  41. private Queue<Message> _sendQueue = new Queue<Message>();
  42.  
  43. private List<Peer> _peers = new List<Peer>();
  44.  
  45. private IPEndPoint _serverEndPoint = null;
  46. private Process _runningProcess = null;
  47.  
  48. public bool Discover { get; set; } = false;
  49. public bool IsServer { get; set; } = false;
  50.  
  51. public void Start()
  52. {
  53. _client = new UdpClient(Port);
  54. _client.EnableBroadcast = true;
  55. Log($"Listening on port {Port}");
  56.  
  57. _listenerTask = Task.Run(ListenerThreadFunc);
  58. _senderTask = Task.Run(SenderThreadFunc);
  59.  
  60. if (Discover)
  61. _discoverTask = Task.Run(DiscoverThreadFunc);
  62. }
  63.  
  64. public EventHandler<Peer> PeerConnected;
  65. public EventHandler<Peer> PeerDisconnected;
  66. public EventHandler<Peer> PeerStatus;
  67.  
  68. private async Task ListenerThreadFunc()
  69. {
  70. while (true)
  71. {
  72. try
  73. {
  74. var result = await _client.ReceiveAsync();
  75. byte id = result.Buffer[0];
  76.  
  77. // Discover request
  78. if (id == 0xC0)
  79. {
  80. //Log($"Received discover from {result.RemoteEndPoint}");
  81. SendComputerInfo(result.RemoteEndPoint);
  82. }
  83.  
  84. // Computer info
  85. else if (id == 0xC1)
  86. {
  87. using (var s = new MemoryStream(result.Buffer))
  88. {
  89. var b = new BinaryReader(s);
  90. b.ReadByte(); // Packet id
  91. var computerName = b.ReadString();
  92. var runningProcess = b.ReadBoolean();
  93.  
  94. lock (_peers)
  95. {
  96. var peer = _peers.FirstOrDefault(p =>
  97. p.EndPoint.Address == result.RemoteEndPoint.Address ||
  98. p.ComputerName == computerName);
  99. if (peer == null)
  100. {
  101. peer = new Peer()
  102. {
  103. EndPoint = result.RemoteEndPoint,
  104. ComputerName = computerName,
  105. LastAlive = DateTime.Now,
  106. };
  107. _peers.Add(peer);
  108. PeerConnected?.Invoke(this, peer);
  109. Log($"{computerName} connected from {peer.EndPoint}");
  110. }
  111. else
  112. {
  113. peer.ComputerName = computerName;
  114. peer.LastAlive = DateTime.Now;
  115. peer.RunningProcess = runningProcess;
  116. PeerStatus?.Invoke(this, peer);
  117. }
  118. }
  119. }
  120. }
  121.  
  122. // Message
  123. else if (id == 0xC2)
  124. {
  125. using (var s = new MemoryStream(result.Buffer))
  126. {
  127. var b = new BinaryReader(s);
  128. b.ReadByte(); // Packet id
  129.  
  130. var message = b.ReadString();
  131. Log($"Message from {result.RemoteEndPoint}: {message}");
  132. }
  133. }
  134.  
  135. // Run process
  136. else if (id == 0xC3 && !IsServer)
  137. {
  138. if (_runningProcess != null)
  139. {
  140. _runningProcess.Refresh();
  141. if (!_runningProcess.HasExited)
  142. continue;
  143. }
  144.  
  145. using (var s = new MemoryStream(result.Buffer))
  146. {
  147. var b = new BinaryReader(s);
  148. b.ReadByte(); // Packet id
  149.  
  150. var filename = b.ReadString();
  151. var arguments = b.ReadString();
  152.  
  153. Log($"{result.RemoteEndPoint} requests process start: {filename} {arguments}");
  154.  
  155. _serverEndPoint = result.RemoteEndPoint;
  156. try
  157. {
  158. _runningProcess = Process.Start(filename,
  159. arguments.Replace("{IP}", _serverEndPoint.Address.ToString()));
  160.  
  161. Thread.Sleep(300);
  162.  
  163. _runningProcess.Refresh();
  164.  
  165. if (_runningProcess.HasExited)
  166. {
  167. SendErrorMessage(_serverEndPoint,
  168. $"Process creation failed with exit code: {_runningProcess.ExitCode}");
  169. Log($"Process creation failed with exit code: {_runningProcess.ExitCode}");
  170. }
  171. }
  172. catch (Exception e)
  173. {
  174. Log($"{e}");
  175. SendErrorMessage(_serverEndPoint, $"Process creation failed: {e}");
  176. }
  177. }
  178. }
  179.  
  180. // Stop process
  181. else if (id == 0xC4)
  182. {
  183. if (_runningProcess != null)
  184. {
  185. _runningProcess.Refresh();
  186. if (_runningProcess.HasExited)
  187. continue;
  188. }
  189. else
  190. continue;
  191.  
  192. try
  193. {
  194. _runningProcess.Kill();
  195. Thread.Sleep(200);
  196.  
  197. _runningProcess.Refresh();
  198.  
  199. if (_runningProcess.HasExited)
  200. SendErrorMessage(_serverEndPoint, $"Process killed successfully");
  201. else
  202. SendErrorMessage(_serverEndPoint, $"Failed to kill the process");
  203. }
  204. catch (Exception e)
  205. {
  206. SendErrorMessage(_serverEndPoint, $"Failed to kill the process: {e}");
  207. }
  208.  
  209. _runningProcess = null;
  210. }
  211.  
  212. // Run with transfer
  213. else if (id == 0xC5 && !IsServer)
  214. {
  215. if (_runningProcess != null)
  216. {
  217. _runningProcess.Refresh();
  218. if (!_runningProcess.HasExited)
  219. continue;
  220. }
  221.  
  222. using (var s = new MemoryStream(result.Buffer))
  223. {
  224. var b = new BinaryReader(s);
  225. b.ReadByte(); // Packet id
  226.  
  227. int port = b.ReadInt32();
  228.  
  229. Thread.Sleep(300);
  230.  
  231. _serverEndPoint = result.RemoteEndPoint;
  232.  
  233. var transferer = new FileTransferer();
  234. transferer.LogMessage += (sender, msg) => Log($"FT: {msg}");
  235. var files = transferer.ReceiveFiles(new IPEndPoint(_serverEndPoint.Address, port));
  236.  
  237. var filename = b.ReadString();
  238. var arguments = b.ReadString();
  239.  
  240. var targetDir = "downloaded";
  241.  
  242. foreach (var file in files)
  243. {
  244. Directory.CreateDirectory(Path.Combine(targetDir, Path.GetDirectoryName(file.Item1)));
  245. File.WriteAllBytes(Path.Combine(targetDir, file.Item1), file.Item2);
  246. }
  247.  
  248. filename = Path.Combine(targetDir, filename);
  249.  
  250. Log($"{result.RemoteEndPoint} requests process start: {filename} {arguments}");
  251.  
  252. _serverEndPoint = result.RemoteEndPoint;
  253. try
  254. {
  255. _runningProcess = Process.Start(filename,
  256. arguments.Replace("{IP}", _serverEndPoint.Address.ToString()));
  257.  
  258. Thread.Sleep(300);
  259.  
  260. _runningProcess.Refresh();
  261.  
  262. if (_runningProcess.HasExited)
  263. {
  264. SendErrorMessage(_serverEndPoint,
  265. $"Process creation failed with exit code: {_runningProcess.ExitCode}");
  266. Log($"Process creation failed with exit code: {_runningProcess.ExitCode}");
  267. }
  268. }
  269. catch (Exception e)
  270. {
  271. Log($"{e}");
  272. SendErrorMessage(_serverEndPoint, $"Process creation failed: {e}");
  273. }
  274. }
  275. }
  276. }
  277. catch (Exception e)
  278. {
  279.  
  280. }
  281. }
  282. }
  283.  
  284. private async Task SenderThreadFunc()
  285. {
  286. while (true)
  287. {
  288. try
  289. {
  290. bool any;
  291.  
  292. lock (_sendQueue)
  293. any = _sendQueue.Any();
  294.  
  295. if (any)
  296. {
  297. Message msg;
  298.  
  299. lock (_sendQueue)
  300. msg = _sendQueue.Dequeue();
  301.  
  302. await _client.SendAsync(msg.Data, msg.Data.Length, msg.EndPoint);
  303. //Log($"Sent data to {msg.EndPoint}");
  304. }
  305. else
  306. await Task.Delay(10);
  307. }
  308. catch (Exception e)
  309. {
  310.  
  311. }
  312. }
  313. }
  314.  
  315. private async Task DiscoverThreadFunc()
  316. {
  317. while (true)
  318. {
  319. try
  320. {
  321. lock (_sendQueue)
  322. {
  323. // Sala 10
  324. _sendQueue.Enqueue(new Message()
  325. {
  326. EndPoint = new IPEndPoint(
  327. IPAddress.Parse("192.168.4.143"),
  328. Port),
  329. Data = new byte[] { 0xC0 },
  330. });
  331.  
  332. foreach (NetworkInterface ni in NetworkInterface.GetAllNetworkInterfaces())
  333. {
  334. if ((ni.NetworkInterfaceType == NetworkInterfaceType.Wireless80211 ||
  335. ni.NetworkInterfaceType == NetworkInterfaceType.Ethernet) &&
  336. ni.OperationalStatus == OperationalStatus.Up)
  337. {
  338. foreach (UnicastIPAddressInformation ip in ni.GetIPProperties().UnicastAddresses)
  339. {
  340. if (ip.Address.AddressFamily == System.Net.Sockets.AddressFamily.InterNetwork)
  341. {
  342. var broadcastIp = ip.Address.GetAddressBytes();
  343. for (int i = 0; i < 4; i++)
  344. if (ip.IPv4Mask.GetAddressBytes()[i] == 0)
  345. broadcastIp[i] |= 0xFF;
  346.  
  347. _sendQueue.Enqueue(new Message()
  348. {
  349. EndPoint = new IPEndPoint(
  350. new IPAddress(broadcastIp),
  351. Port),
  352. Data = new byte[] { 0xC0 },
  353. });
  354. }
  355. }
  356. }
  357. }
  358. }
  359.  
  360. // Check timeouts
  361. Peer[] peers;
  362. lock (_peers)
  363. peers = _peers.ToArray();
  364.  
  365. foreach (var peer in peers)
  366. {
  367. if (DateTime.Now > peer.LastAlive.AddMilliseconds(Timeout))
  368. {
  369. _peers.Remove(peer);
  370. PeerDisconnected?.Invoke(this, peer);
  371. Log($"{peer.EndPoint} disconnected");
  372. }
  373. }
  374. }
  375. catch (Exception e)
  376. {
  377.  
  378. }
  379.  
  380. await Task.Delay(1000);
  381. }
  382. }
  383.  
  384. public EventHandler<string> LogMessage;
  385. private void Log(string message)
  386. {
  387. LogMessage?.Invoke(this, message);
  388. }
  389.  
  390. private void SendComputerInfo(IPEndPoint endPoint)
  391. {
  392. var s = new MemoryStream();
  393. var b = new BinaryWriter(s);
  394.  
  395. b.Write((byte)0xC1);
  396. b.Write(Environment.MachineName);
  397. if (_runningProcess != null)
  398. _runningProcess.Refresh();
  399. b.Write(!_runningProcess?.HasExited ?? false);
  400.  
  401. SendMessage(new Message()
  402. {
  403. EndPoint = endPoint,
  404. Data = s.ToArray(),
  405. });
  406.  
  407. //Log($"Sending computer info to {endPoint}");
  408. }
  409.  
  410. private void SendMessage(Message msg)
  411. {
  412. lock (_sendQueue)
  413. _sendQueue.Enqueue(msg);
  414. }
  415.  
  416. private void SendErrorMessage(IPEndPoint endPoint, string message)
  417. {
  418. var s = new MemoryStream();
  419. var b = new BinaryWriter(s);
  420.  
  421. b.Write((byte)0xC2);
  422. b.Write(message);
  423.  
  424. SendMessage(new Message()
  425. {
  426. EndPoint = endPoint,
  427. Data = s.ToArray(),
  428. });
  429. }
  430.  
  431. public void SendProcessStart(string filename, string args)
  432. {
  433. var r = new Random();
  434.  
  435. var s = new MemoryStream();
  436. var b = new BinaryWriter(s);
  437.  
  438. b.Write((byte)0xC3);
  439. b.Write(filename);
  440. b.Write(args.Replace("{RANDOM}", "" + r.Next(10000, 99999)));
  441.  
  442. lock (_peers)
  443. foreach (var peer in _peers)
  444. {
  445. SendMessage(new Message()
  446. {
  447. EndPoint = peer.EndPoint,
  448. Data = s.ToArray(),
  449. });
  450. }
  451. }
  452.  
  453. public async Task SendProcessStartWithFiles(string filename, string args)
  454. {
  455. var r = new Random();
  456.  
  457. var transferer = new FileTransferer();
  458. transferer.LogMessage += (sender, msg) => Log($"FT: {msg}");
  459.  
  460. try
  461. {
  462. var mainDir = Path.GetDirectoryName(filename);
  463.  
  464. var files = new List<string>();
  465.  
  466. foreach (var file in Directory.GetFileSystemEntries(mainDir, "*", SearchOption.AllDirectories))
  467. if (!Directory.Exists(file))
  468. files.Add(file);
  469.  
  470. transferer.Start(files.ToArray());
  471. }
  472. catch (Exception e)
  473. {
  474. MessageBox.Show($"Failed to send files: {e}", "ClusterX", MessageBoxButton.OK, MessageBoxImage.Error);
  475. return;
  476. }
  477.  
  478. Thread.Sleep(100);
  479.  
  480. var s = new MemoryStream();
  481. var b = new BinaryWriter(s);
  482.  
  483. b.Write((byte)0xC5);
  484. b.Write(transferer.Port);
  485. b.Write(filename);
  486. b.Write(args.Replace("{RANDOM}", "" + r.Next(10000, 99999)));
  487.  
  488. lock (_peers)
  489. foreach (var peer in _peers)
  490. {
  491. SendMessage(new Message()
  492. {
  493. EndPoint = peer.EndPoint,
  494. Data = s.ToArray(),
  495. });
  496. }
  497. }
  498.  
  499. public async Task SendProcessKill()
  500. {
  501. var s = new MemoryStream();
  502. var b = new BinaryWriter(s);
  503.  
  504. b.Write((byte)0xC4);
  505.  
  506. lock (_peers)
  507. foreach (var peer in _peers)
  508. {
  509. SendMessage(new Message()
  510. {
  511. EndPoint = peer.EndPoint,
  512. Data = s.ToArray(),
  513. });
  514. }
  515. }
  516. }
  517. }
Advertisement
Add Comment
Please, Sign In to add comment