RosBridgeClient.h 16.2 KB
Newer Older
Valentin Platzgummer's avatar
Valentin Platzgummer committed
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16
#pragma once

/*
*  Created on: Apr 16, 2018
*      Author: Poom Pianpak
*/

#include "Simple-WebSocket-Server/client_ws.hpp"

#include "rapidjson/include/rapidjson/document.h"
#include "rapidjson/include/rapidjson/writer.h"
#include "rapidjson/include/rapidjson/stringbuffer.h"

#include <chrono>
#include <functional>
#include <thread>
17
#include <future>
18
#include <mutex>
Valentin Platzgummer's avatar
Valentin Platzgummer committed
19 20 21 22 23 24

using WsClient = SimpleWeb::SocketClient<SimpleWeb::WS>;
using InMessage =  std::function<void(std::shared_ptr<WsClient::Connection>, std::shared_ptr<WsClient::InMessage>)>;

class RosbridgeWsClient
{
25 26 27 28 29 30 31 32 33 34
  std::string server_port_path;
  std::unordered_map<std::string, std::shared_ptr<WsClient>> client_map;
  std::mutex map_mutex;

  void start(const std::string& client_name, std::shared_ptr<WsClient> client, const std::string& message)
  {
    if (!client->on_open)
    {
#ifdef DEBUG
      client->on_open = [client_name, message](std::shared_ptr<WsClient::Connection> connection) {
Valentin Platzgummer's avatar
Valentin Platzgummer committed
35
#else
36
      client->on_open = [message](std::shared_ptr<WsClient::Connection> connection) {
Valentin Platzgummer's avatar
Valentin Platzgummer committed
37 38
#endif

39 40 41
#ifdef DEBUG
        std::cout << client_name << ": Opened connection" << std::endl;
        std::cout << client_name << ": Sending message: " << message << std::endl;
Valentin Platzgummer's avatar
Valentin Platzgummer committed
42
#endif
43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68
        connection->send(message);
      };
    }

#ifdef DEBUG
    if (!client->on_message)
    {
      client->on_message = [client_name](std::shared_ptr<WsClient::Connection> /*connection*/, std::shared_ptr<WsClient::InMessage> in_message) {
        std::cout << client_name << ": Message received: " << in_message->string() << std::endl;
      };
    }

    if (!client->on_close)
    {
      client->on_close = [client_name](std::shared_ptr<WsClient::Connection> /*connection*/, int status, const std::string & /*reason*/) {
        std::cout << client_name << ": Closed connection with status code " << status << std::endl;
      };
    }

    if (!client->on_error)
    {
      // See http://www.boost.org/doc/libs/1_55_0/doc/html/boost_asio/reference.html, Error Codes for error code meanings
      client->on_error = [client_name](std::shared_ptr<WsClient::Connection> /*connection*/, const SimpleWeb::error_code &ec) {
        std::cout << client_name << ": Error: " << ec << ", error message: " << ec.message() << std::endl;
      };
    }
Valentin Platzgummer's avatar
Valentin Platzgummer committed
69 70
#endif

71 72
#ifdef DEBUG
    std::thread client_thread([client_name, client]() {
Valentin Platzgummer's avatar
Valentin Platzgummer committed
73
#else
74
    std::thread client_thread([client]() {
Valentin Platzgummer's avatar
Valentin Platzgummer committed
75
#endif
76
      client->start();
Valentin Platzgummer's avatar
Valentin Platzgummer committed
77

78 79
#ifdef DEBUG
      std::cout << client_name << ": Terminated" << std::endl;
Valentin Platzgummer's avatar
Valentin Platzgummer committed
80
#endif
81 82 83 84 85
      client->on_open = NULL;
      client->on_message = NULL;
      client->on_close = NULL;
      client->on_error = NULL;
    });
Valentin Platzgummer's avatar
Valentin Platzgummer committed
86

87 88
    client_thread.detach();
  }
Valentin Platzgummer's avatar
Valentin Platzgummer committed
89 90

public:
91 92 93 94 95 96 97 98 99 100 101 102 103 104
  RosbridgeWsClient(const std::string& server_port_path)
  {
    this->server_port_path = server_port_path;
  }

  ~RosbridgeWsClient()
  {
    std::lock_guard<std::mutex> lk(map_mutex); // neccessary?
    for (auto& client : client_map)
    {
      client.second->stop();
      client.second.reset();
    }
  }
Valentin Platzgummer's avatar
Valentin Platzgummer committed
105

106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132
 // The execution can take up to 100 ms!
 bool connected(){

    auto p = std::make_shared<std::promise<void>>();
    auto future = p->get_future();

    auto callback = [](std::shared_ptr<WsClient::Connection> connection, std::shared_ptr<std::promise<void>> p) {
        p->set_value();
        connection->send_close(1000);
    };
    std::shared_ptr<WsClient> status_client = std::make_shared<WsClient>(server_port_path);
    status_client->on_open = std::bind(callback, std::placeholders::_1, p);

    std::async([status_client]{

        status_client->start();
        status_client->on_open = NULL;
        status_client->on_message = NULL;
        status_client->on_close = NULL;
        status_client->on_error = NULL;

    });

    auto status = future.wait_for(std::chrono::milliseconds(100));
    return status == std::future_status::ready;
 }

133 134 135 136 137 138 139 140 141 142 143 144 145 146
  // Adds a client to the client_map.
  void addClient(const std::string& client_name)
  {
    std::lock_guard<std::mutex> lk(map_mutex);
    std::unordered_map<std::string, std::shared_ptr<WsClient>>::iterator it = client_map.find(client_name);
    if (it == client_map.end())
    {
      client_map[client_name] = std::make_shared<WsClient>(server_port_path);
    }
#ifdef DEBUG
    else
    {
      std::cerr << client_name << " has already been created" << std::endl;
    }
Valentin Platzgummer's avatar
Valentin Platzgummer committed
147
#endif
148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173
  }

  std::shared_ptr<WsClient> getClient(const std::string& client_name)
  {
    std::lock_guard<std::mutex> lk(map_mutex);
    std::unordered_map<std::string, std::shared_ptr<WsClient>>::iterator it = client_map.find(client_name);
    if (it != client_map.end())
    {
      return it->second;
    }
    return NULL;
  }

  void stopClient(const std::string& client_name)
  {
    std::lock_guard<std::mutex> lk(map_mutex);
    std::unordered_map<std::string, std::shared_ptr<WsClient>>::iterator it = client_map.find(client_name);
    if (it != client_map.end())
    {
      it->second->stop();
      it->second->on_open = NULL;
      it->second->on_message = NULL;
      it->second->on_close = NULL;
      it->second->on_error = NULL;
#ifdef DEBUG
      std::cout << client_name << " has been stopped" << std::endl;
Valentin Platzgummer's avatar
Valentin Platzgummer committed
174
#endif
175 176 177 178 179 180
    }
#ifdef DEBUG
    else
    {
      std::cerr << client_name << " has not been created" << std::endl;
    }
Valentin Platzgummer's avatar
Valentin Platzgummer committed
181
#endif
182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201
  }

  void removeClient(const std::string& client_name)
  {
    std::lock_guard<std::mutex> lk(map_mutex);
    std::unordered_map<std::string, std::shared_ptr<WsClient>>::iterator it = client_map.find(client_name);
    if (it != client_map.end())
    {
      // Stop the client asynchronously in 100 ms.
      // This is to ensure, that all threads involving the client have been launched.
      std::thread t([](std::shared_ptr<WsClient> client){
         std::this_thread::sleep_for(std::chrono::milliseconds(100));
         client->stop();
         client.reset();
      },
      it->second /*lambda param*/ );
      client_map.erase(it);
      t.detach();
#ifdef DEBUG
      std::cout << client_name << " has been removed" << std::endl;
Valentin Platzgummer's avatar
Valentin Platzgummer committed
202
#endif
203 204 205 206 207 208
    }
#ifdef DEBUG
    else
    {
      std::cerr << client_name << " has not been created" << std::endl;
    }
Valentin Platzgummer's avatar
Valentin Platzgummer committed
209
#endif
210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234
  }

  // Gets the client from client_map and starts it. Advertising is essentially sending a message.
  // One client per topic!
  void advertise(const std::string& client_name, const std::string& topic, const std::string& type, const std::string& id = "")
  {
    std::lock_guard<std::mutex> lk(map_mutex);
    std::unordered_map<std::string, std::shared_ptr<WsClient>>::iterator it = client_map.find(client_name);
    if (it != client_map.end())
    {
      std::string message = "\"op\":\"advertise\", \"topic\":\"" + topic + "\", \"type\":\"" + type + "\"";

      if (id.compare("") != 0)
      {
        message += ", \"id\":\"" + id + "\"";
      }
      message = "{" + message + "}";

      start(client_name, it->second, message);
    }
#ifdef DEBUG
    else
    {
      std::cerr << client_name << "has not been created" << std::endl;
    }
Valentin Platzgummer's avatar
Valentin Platzgummer committed
235
#endif
236
  }
Valentin Platzgummer's avatar
Valentin Platzgummer committed
237

238 239 240 241 242
  void publish(const std::string& topic, const rapidjson::Document& msg, const std::string& id = "")
  {
    rapidjson::StringBuffer strbuf;
    rapidjson::Writer<rapidjson::StringBuffer> writer(strbuf);
    msg.Accept(writer);
Valentin Platzgummer's avatar
Valentin Platzgummer committed
243

244
    std::string message = "\"op\":\"publish\", \"topic\":\"" + topic + "\", \"msg\":" + strbuf.GetString();
Valentin Platzgummer's avatar
Valentin Platzgummer committed
245

246 247 248 249 250
    if (id.compare("") != 0)
    {
      message += ", \"id\":\"" + id + "\"";
    }
    message = "{" + message + "}";
Valentin Platzgummer's avatar
Valentin Platzgummer committed
251

252
    std::shared_ptr<WsClient> publish_client = std::make_shared<WsClient>(server_port_path);
Valentin Platzgummer's avatar
Valentin Platzgummer committed
253

254 255 256 257
    publish_client->on_open = [message](std::shared_ptr<WsClient::Connection> connection) {
#ifdef DEBUG
      std::cout << "publish_client: Opened connection" << std::endl;
      std::cout << "publish_client: Sending message: " << message << std::endl;
Valentin Platzgummer's avatar
Valentin Platzgummer committed
258
#endif
259
      connection->send(message);
Valentin Platzgummer's avatar
Valentin Platzgummer committed
260

261 262 263
      // TODO: This could be improved by creating a thread to keep publishing the message instead of closing it right away
      connection->send_close(1000);
    };
Valentin Platzgummer's avatar
Valentin Platzgummer committed
264

265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309
    start("publish_client", publish_client, message);
  }

  void subscribe(const std::string& client_name, const std::string& topic, const InMessage& callback, const std::string& id = "", const std::string& type = "", int throttle_rate = -1, int queue_length = -1, int fragment_size = -1, const std::string& compression = "")
  {
    std::lock_guard<std::mutex> lk(map_mutex);
    std::unordered_map<std::string, std::shared_ptr<WsClient>>::iterator it = client_map.find(client_name);
    if (it != client_map.end())
    {
      std::string message = "\"op\":\"subscribe\", \"topic\":\"" + topic + "\"";

      if (id.compare("") != 0)
      {
        message += ", \"id\":\"" + id + "\"";
      }
      if (type.compare("") != 0)
      {
        message += ", \"type\":\"" + type + "\"";
      }
      if (throttle_rate > -1)
      {
        message += ", \"throttle_rate\":" + std::to_string(throttle_rate);
      }
      if (queue_length > -1)
      {
        message += ", \"queue_length\":" + std::to_string(queue_length);
      }
      if (fragment_size > -1)
      {
        message += ", \"fragment_size\":" + std::to_string(fragment_size);
      }
      if (compression.compare("none") == 0 || compression.compare("png") == 0)
      {
        message += ", \"compression\":\"" + compression + "\"";
      }
      message = "{" + message + "}";

      it->second->on_message = callback;
      start(client_name, it->second, message);
    }
#ifdef DEBUG
    else
    {
      std::cerr << client_name << "has not been created" << std::endl;
    }
Valentin Platzgummer's avatar
Valentin Platzgummer committed
310
#endif
311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328
  }

  void advertiseService(const std::string& client_name, const std::string& service, const std::string& type, const InMessage& callback)
  {
    std::lock_guard<std::mutex> lk(map_mutex);
    std::unordered_map<std::string, std::shared_ptr<WsClient>>::iterator it = client_map.find(client_name);
    if (it != client_map.end())
    {
      std::string message = "{\"op\":\"advertise_service\", \"service\":\"" + service + "\", \"type\":\"" + type + "\"}";

      it->second->on_message = callback;
      start(client_name, it->second, message);
    }
#ifdef DEBUG
    else
    {
      std::cerr << client_name << "has not been created" << std::endl;
    }
Valentin Platzgummer's avatar
Valentin Platzgummer committed
329
#endif
330
  }
Valentin Platzgummer's avatar
Valentin Platzgummer committed
331

332
  void serviceResponse(const std::string& service, const std::string& id, bool result, const rapidjson::Document& values)
333 334
  {
    std::string message = "\"op\":\"service_response\", \"service\":\"" + service + "\", \"result\":" + ((result)? "true" : "false");
Valentin Platzgummer's avatar
Valentin Platzgummer committed
335

336 337 338
    // Rosbridge somehow does not allow service_response to be sent without id and values
    // , so we cannot omit them even though the documentation says they are optional.
    message += ", \"id\":\"" + id + "\"";
Valentin Platzgummer's avatar
Valentin Platzgummer committed
339

340 341 342 343
    // Convert JSON document to string
    rapidjson::StringBuffer strbuf;
    rapidjson::Writer<rapidjson::StringBuffer> writer(strbuf);
    values.Accept(writer);
Valentin Platzgummer's avatar
Valentin Platzgummer committed
344

345 346
    message += ", \"values\":" + std::string(strbuf.GetString());
    message = "{" + message + "}";
Valentin Platzgummer's avatar
Valentin Platzgummer committed
347

348
    std::shared_ptr<WsClient> service_response_client = std::make_shared<WsClient>(server_port_path);
Valentin Platzgummer's avatar
Valentin Platzgummer committed
349

350 351 352 353
    service_response_client->on_open = [message](std::shared_ptr<WsClient::Connection> connection) {
#ifdef DEBUG
      std::cout << "service_response_client: Opened connection" << std::endl;
      std::cout << "service_response_client: Sending message: " << message << std::endl;
Valentin Platzgummer's avatar
Valentin Platzgummer committed
354
#endif
355
      connection->send(message);
Valentin Platzgummer's avatar
Valentin Platzgummer committed
356

357 358
      connection->send_close(1000);
    };
Valentin Platzgummer's avatar
Valentin Platzgummer committed
359

360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400
    start("service_response_client", service_response_client, message);
  }

  void callService(const std::string& service, const InMessage& callback, const rapidjson::Document& args = {}, const std::string& id = "", int fragment_size = -1, const std::string& compression = "")
  {
    std::string message = "\"op\":\"call_service\", \"service\":\"" + service + "\"";

    if (!args.IsNull())
    {
      rapidjson::StringBuffer strbuf;
      rapidjson::Writer<rapidjson::StringBuffer> writer(strbuf);
      args.Accept(writer);

      message += ", \"args\":" + std::string(strbuf.GetString());
    }
    if (id.compare("") != 0)
    {
      message += ", \"id\":\"" + id + "\"";
    }
    if (fragment_size > -1)
    {
      message += ", \"fragment_size\":" + std::to_string(fragment_size);
    }
    if (compression.compare("none") == 0 || compression.compare("png") == 0)
    {
      message += ", \"compression\":\"" + compression + "\"";
    }
    message = "{" + message + "}";

    std::shared_ptr<WsClient> call_service_client = std::make_shared<WsClient>(server_port_path);

    if (callback)
    {
      call_service_client->on_message = callback;
    }
    else
    {
      call_service_client->on_message = [](std::shared_ptr<WsClient::Connection> connection, std::shared_ptr<WsClient::InMessage> in_message) {
#ifdef DEBUG
        std::cout << "call_service_client: Message received: " << in_message->string() << std::endl;
        std::cout << "call_service_client: Sending close connection" << std::endl;
Valentin Platzgummer's avatar
Valentin Platzgummer committed
401
#endif
402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514
        connection->send_close(1000);
      };
    }

    start("call_service_client", call_service_client, message);
  }

  //!
  //! \brief Checks if the service \p service is available.
  //! \param service Service name.
  //! \return Returns true if the service is available, false either.
  //! \note Don't call this function to frequently. Use \fn waitForService() instead.
  //!
  bool serviceAvailable(const std::string& service)
  {
    const rapidjson::Document args = {};
    const std::string& id = "";
    const int fragment_size = -1;
    const std::string& compression = "";
    std::string message = "\"op\":\"call_service\", \"service\":\"" + service + "\"";
    std::string client_name("wait_for_service_client");

    if (!args.IsNull())
    {
      rapidjson::StringBuffer strbuf;
      rapidjson::Writer<rapidjson::StringBuffer> writer(strbuf);
      args.Accept(writer);

      message += ", \"args\":" + std::string(strbuf.GetString());
    }
    if (id.compare("") != 0)
    {
      message += ", \"id\":\"" + id + "\"";
    }
    if (fragment_size > -1)
    {
      message += ", \"fragment_size\":" + std::to_string(fragment_size);
    }
    if (compression.compare("none") == 0 || compression.compare("png") == 0)
    {
      message += ", \"compression\":\"" + compression + "\"";
    }
    message = "{" + message + "}";

    std::shared_ptr<WsClient> wait_for_service_client = std::make_shared<WsClient>(server_port_path);
    std::shared_ptr<std::promise<bool>> p(std::make_shared<std::promise<bool>>());
    auto future = p->get_future();
#ifdef DEBUG
    wait_for_service_client->on_message = [p, service, client_name](
#else
    wait_for_service_client->on_message = [p](
#endif
            std::shared_ptr<WsClient::Connection> connection,
            std::shared_ptr<WsClient::InMessage> in_message){
#ifdef DEBUG
#endif
        rapidjson::Document doc;
        doc.Parse(in_message->string().c_str());
        if ( !doc.HasParseError()
             && doc.HasMember("result")
             && doc["result"].IsBool()
             && doc["result"] == true )
        {
#ifdef DEBUG
            std::cout << client_name << ": "
                      << "service " << service
                      << " available." << std::endl;
            std::cout << client_name << ": "
                      << "message: " << in_message->string()
                      << std::endl;
#endif
            p->set_value(true);
        } else {
            p->set_value(false);
        }
        connection->send_close(1000);
    };
    start(client_name, wait_for_service_client, message);
    return future.get();
  }

  //!
  //! \brief Waits until the service with the name \p service is available.
  //! \param service Service name.
  //! @note This function will block as long as the service is not available.
  //!
  void waitForService(const std::string& service)
  {


#ifdef DEBUG
    auto s = std::chrono::high_resolution_clock::now();
    long counter = 0;
#endif
    while(true)
    {
#ifdef DEBUG
        ++counter;
#endif
        if (serviceAvailable(service)){
            break;
        }
        std::this_thread::sleep_for(std::chrono::milliseconds(50)); // Avoid excessive polling.
    };
#ifdef DEBUG
    auto e = std::chrono::high_resolution_clock::now();
    std::cout << "waitForService(): time: "
              << std::chrono::duration_cast<std::chrono::milliseconds>(e-s).count()
              << " ms." << std::endl;
    std::cout << "waitForService(): clients launched: "
              << counter << std::endl;
#endif
  }
Valentin Platzgummer's avatar
Valentin Platzgummer committed
515
};