2018-09-22 23:11:22 +02:00
|
|
|
#include <memory>
|
|
|
|
|
|
2009-11-24 18:36:45 +00:00
|
|
|
// Cyphesis Online RPG Server and AI Engine
|
|
|
|
|
// Copyright (C) 2009 Alistair Riddoch
|
|
|
|
|
//
|
|
|
|
|
// This program is free software; you can redistribute it and/or modify
|
|
|
|
|
// it under the terms of the GNU General Public License as published by
|
|
|
|
|
// the Free Software Foundation; either version 2 of the License, or
|
|
|
|
|
// (at your option) any later version.
|
|
|
|
|
//
|
|
|
|
|
// This program is distributed in the hope that it will be useful,
|
|
|
|
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
|
|
|
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
|
|
|
|
// GNU General Public License for more details.
|
|
|
|
|
//
|
|
|
|
|
// You should have received a copy of the GNU General Public License
|
|
|
|
|
// along with this program; if not, write to the Free Software Foundation,
|
|
|
|
|
// Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
|
|
|
|
|
|
|
|
|
|
2009-11-25 12:36:36 +00:00
|
|
|
#ifdef HAVE_CONFIG_H
|
|
|
|
|
#endif // HAVE_CONFIG_H
|
|
|
|
|
|
2009-11-24 18:36:45 +00:00
|
|
|
#include "common/AtlasStreamClient.h"
|
2009-11-28 18:31:28 +00:00
|
|
|
#include "common/ClientTask.h"
|
2009-11-24 18:36:45 +00:00
|
|
|
|
2012-05-29 12:29:29 +01:00
|
|
|
#include "common/debug.h"
|
2018-09-22 23:11:22 +02:00
|
|
|
#include "AtlasStreamClient.h"
|
|
|
|
|
|
2012-05-29 12:29:29 +01:00
|
|
|
|
2009-11-24 18:36:45 +00:00
|
|
|
#include <Atlas/Codec.h>
|
|
|
|
|
#include <Atlas/Objects/Anonymous.h>
|
|
|
|
|
#include <Atlas/Objects/Encoder.h>
|
|
|
|
|
#include <Atlas/Net/Stream.h>
|
|
|
|
|
|
2010-07-10 00:06:00 +01:00
|
|
|
#ifdef _WIN32
|
|
|
|
|
#undef DATADIR
|
|
|
|
|
#endif // _WIN32
|
|
|
|
|
|
2014-02-02 01:05:39 +01:00
|
|
|
#include <iostream>
|
2009-11-29 20:40:50 +00:00
|
|
|
|
2009-11-25 05:08:36 +00:00
|
|
|
using Atlas::Message::Element;
|
|
|
|
|
using Atlas::Message::ListType;
|
|
|
|
|
using Atlas::Message::MapType;
|
2009-11-27 14:39:42 +00:00
|
|
|
using Atlas::Objects::Root;
|
2009-11-24 18:36:45 +00:00
|
|
|
using Atlas::Objects::Entity::Anonymous;
|
|
|
|
|
using Atlas::Objects::Operation::Create;
|
|
|
|
|
using Atlas::Objects::Operation::Login;
|
|
|
|
|
using Atlas::Objects::Operation::RootOperation;
|
|
|
|
|
|
2014-02-02 01:01:53 +01:00
|
|
|
using namespace boost::asio;
|
|
|
|
|
|
2018-10-04 21:50:03 +02:00
|
|
|
static const bool debug_flag = false;
|
2016-07-09 23:30:05 +02:00
|
|
|
|
2014-02-02 01:01:53 +01:00
|
|
|
|
2014-02-04 22:24:27 +01:00
|
|
|
StreamClientSocketBase::StreamClientSocketBase(boost::asio::io_service& io_service, std::function<void()>& dispatcher)
|
2018-09-22 23:11:22 +02:00
|
|
|
: m_io_service(io_service),
|
|
|
|
|
mDispatcher(dispatcher),
|
|
|
|
|
m_ios(&mBuffer),
|
|
|
|
|
m_codec(nullptr),
|
|
|
|
|
m_encoder(nullptr),
|
|
|
|
|
m_is_connected(false)
|
2014-02-02 01:01:53 +01:00
|
|
|
{
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
StreamClientSocketBase::~StreamClientSocketBase()
|
|
|
|
|
{
|
|
|
|
|
delete m_encoder;
|
|
|
|
|
delete m_codec;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
std::iostream& StreamClientSocketBase::getIos()
|
|
|
|
|
{
|
|
|
|
|
return m_ios;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
int StreamClientSocketBase::negotiate(Atlas::Objects::ObjectsDecoder& decoder)
|
|
|
|
|
{
|
|
|
|
|
|
2016-11-12 23:44:56 +01:00
|
|
|
Atlas::Net::StreamConnect conn("cyphesis_client", m_ios, m_ios);
|
2014-02-02 01:01:53 +01:00
|
|
|
|
|
|
|
|
while (conn.getState() == Atlas::Net::StreamConnect::IN_PROGRESS) {
|
|
|
|
|
write();
|
2014-02-06 13:38:32 +01:00
|
|
|
auto dataReceived = read_blocking();
|
2014-02-02 01:01:53 +01:00
|
|
|
conn.poll(dataReceived > 0);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if (conn.getState() == Atlas::Net::StreamConnect::FAILED) {
|
|
|
|
|
std::cerr << "Failed to negotiate" << std::endl;
|
|
|
|
|
return -1;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
m_codec = conn.getCodec(decoder);
|
|
|
|
|
|
|
|
|
|
if (!m_codec) {
|
|
|
|
|
return -1;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
m_encoder = new Atlas::Objects::ObjectsEncoder(*m_codec);
|
|
|
|
|
|
|
|
|
|
m_codec->streamBegin();
|
|
|
|
|
|
|
|
|
|
do_read();
|
|
|
|
|
|
|
|
|
|
return 0;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
Atlas::Codec& StreamClientSocketBase::getCodec()
|
|
|
|
|
{
|
|
|
|
|
return *m_codec;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
Atlas::Objects::ObjectsEncoder& StreamClientSocketBase::getEncoder()
|
|
|
|
|
{
|
|
|
|
|
return *m_encoder;
|
|
|
|
|
}
|
|
|
|
|
|
2015-05-05 22:01:11 +02:00
|
|
|
int StreamClientSocketBase::poll(const boost::posix_time::time_duration& duration) {
|
|
|
|
|
return poll(duration, []()->bool{return false;});
|
|
|
|
|
}
|
|
|
|
|
|
2018-05-01 22:16:38 +02:00
|
|
|
int StreamClientSocketBase::poll(const boost::posix_time::time_duration& duration, const std::function<bool()>& exitCheckerFn)
|
2014-02-02 01:01:53 +01:00
|
|
|
{
|
|
|
|
|
bool hasExpired = false;
|
2014-10-01 22:42:19 +02:00
|
|
|
bool isCancelled = false;
|
2014-02-03 22:00:46 +01:00
|
|
|
deadline_timer timer(m_io_service);
|
2014-10-01 22:42:19 +02:00
|
|
|
timer.expires_from_now(duration);
|
2014-02-02 01:01:53 +01:00
|
|
|
timer.async_wait([&](boost::system::error_code ec){
|
|
|
|
|
if (!ec) {
|
|
|
|
|
hasExpired = true;
|
2014-10-01 22:42:19 +02:00
|
|
|
} else {
|
|
|
|
|
isCancelled = true;
|
2014-02-02 01:01:53 +01:00
|
|
|
}
|
|
|
|
|
});
|
2014-10-01 22:42:19 +02:00
|
|
|
|
|
|
|
|
//We'll try to only run one handler each polling. Either our timer gets called, or one of the network handlers.
|
|
|
|
|
//The reason for this loop is that when we cancel the timer we need to poll run handlers until the timer handler
|
|
|
|
|
//has been run, since it references locally scoped variables.
|
2015-05-05 22:01:11 +02:00
|
|
|
while (!hasExpired && !isCancelled && !exitCheckerFn()) {
|
2014-10-01 22:42:19 +02:00
|
|
|
m_io_service.run_one();
|
|
|
|
|
//Check if we didn't run the timer handler; if so we should cancel it and then keep on polling until
|
|
|
|
|
//it's been run.
|
|
|
|
|
if (!hasExpired && !isCancelled) {
|
|
|
|
|
timer.cancel();
|
|
|
|
|
}
|
|
|
|
|
}
|
2014-02-02 01:01:53 +01:00
|
|
|
if (!m_is_connected) {
|
|
|
|
|
return -1;
|
|
|
|
|
}
|
|
|
|
|
if (hasExpired) {
|
|
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
return 0;
|
|
|
|
|
}
|
|
|
|
|
|
2018-09-22 23:11:22 +02:00
|
|
|
bool StreamClientSocketBase::isConnected() const
|
|
|
|
|
{
|
|
|
|
|
return m_is_connected;
|
|
|
|
|
}
|
2014-02-02 01:01:53 +01:00
|
|
|
|
|
|
|
|
|
2014-02-04 22:24:27 +01:00
|
|
|
TcpStreamClientSocket::TcpStreamClientSocket(boost::asio::io_service& io_service, std::function<void()>& dispatcher, boost::asio::ip::tcp::endpoint endpoint)
|
|
|
|
|
: StreamClientSocketBase(io_service, dispatcher), m_socket(io_service)
|
2014-02-02 01:01:53 +01:00
|
|
|
{
|
|
|
|
|
m_socket.connect(endpoint);
|
|
|
|
|
m_is_connected = true;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void TcpStreamClientSocket::do_read()
|
|
|
|
|
{
|
|
|
|
|
m_socket.async_read_some(mReadBuffer.prepare(read_buffer_size),
|
|
|
|
|
[this](boost::system::error_code ec, std::size_t length)
|
|
|
|
|
{
|
|
|
|
|
if (!ec)
|
|
|
|
|
{
|
|
|
|
|
mReadBuffer.commit(length);
|
|
|
|
|
this->m_ios.rdbuf(&mReadBuffer);
|
2018-01-28 21:23:18 +01:00
|
|
|
m_codec->poll(true);
|
2014-02-02 01:01:53 +01:00
|
|
|
this->m_ios.rdbuf(&mBuffer);
|
2014-02-04 22:24:27 +01:00
|
|
|
mDispatcher();
|
2014-02-02 01:01:53 +01:00
|
|
|
this->do_read();
|
2014-09-10 19:59:08 +02:00
|
|
|
} else {
|
|
|
|
|
m_is_connected = false;
|
2014-02-02 01:01:53 +01:00
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
|
2017-05-27 20:27:48 +02:00
|
|
|
size_t TcpStreamClientSocket::read_blocking()
|
2014-02-02 01:01:53 +01:00
|
|
|
{
|
|
|
|
|
if (!m_socket.is_open()) {
|
2017-05-27 20:27:48 +02:00
|
|
|
return 0;
|
2014-02-02 01:01:53 +01:00
|
|
|
}
|
|
|
|
|
if (m_socket.available() > 0) {
|
|
|
|
|
auto received = m_socket.read_some(mBuffer.prepare(m_socket.available()));
|
|
|
|
|
if (received > 0) {
|
|
|
|
|
mBuffer.commit(received);
|
|
|
|
|
}
|
|
|
|
|
return received;
|
|
|
|
|
}
|
|
|
|
|
return 0;
|
|
|
|
|
}
|
|
|
|
|
|
2017-05-27 20:27:48 +02:00
|
|
|
size_t TcpStreamClientSocket::write()
|
2014-02-02 01:01:53 +01:00
|
|
|
{
|
|
|
|
|
if (!m_socket.is_open()) {
|
2017-05-27 20:27:48 +02:00
|
|
|
return 0;
|
2014-02-02 01:01:53 +01:00
|
|
|
}
|
|
|
|
|
if (mBuffer.size() == 0) {
|
|
|
|
|
return 0;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
auto size = boost::asio::write(m_socket, mBuffer.data(),
|
|
|
|
|
boost::asio::transfer_all());
|
|
|
|
|
mBuffer.consume(size);
|
|
|
|
|
return size;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
2014-02-04 22:24:27 +01:00
|
|
|
LocalStreamClientSocket::LocalStreamClientSocket(boost::asio::io_service& io_service, std::function<void()>& dispatcher, boost::asio::local::stream_protocol::endpoint endpoint)
|
|
|
|
|
: StreamClientSocketBase(io_service, dispatcher), m_socket(io_service)
|
2014-02-02 01:01:53 +01:00
|
|
|
{
|
|
|
|
|
m_socket.connect(endpoint);
|
|
|
|
|
m_is_connected = true;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void LocalStreamClientSocket::do_read()
|
|
|
|
|
{
|
|
|
|
|
m_socket.async_read_some(mReadBuffer.prepare(read_buffer_size),
|
|
|
|
|
[this](boost::system::error_code ec, std::size_t length)
|
|
|
|
|
{
|
|
|
|
|
if (!ec)
|
|
|
|
|
{
|
|
|
|
|
mReadBuffer.commit(length);
|
|
|
|
|
this->m_ios.rdbuf(&mReadBuffer);
|
2018-01-28 21:23:18 +01:00
|
|
|
m_codec->poll(true);
|
2014-02-02 01:01:53 +01:00
|
|
|
this->m_ios.rdbuf(&mBuffer);
|
2014-02-04 22:24:27 +01:00
|
|
|
mDispatcher();
|
2014-02-02 01:01:53 +01:00
|
|
|
this->do_read();
|
2014-09-10 19:59:08 +02:00
|
|
|
} else {
|
|
|
|
|
m_is_connected = false;
|
2014-02-02 01:01:53 +01:00
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
|
2017-05-27 20:27:48 +02:00
|
|
|
size_t LocalStreamClientSocket::read_blocking()
|
2014-02-02 01:01:53 +01:00
|
|
|
{
|
|
|
|
|
if (!m_socket.is_open()) {
|
2017-05-27 20:27:48 +02:00
|
|
|
return 0;
|
2014-02-02 01:01:53 +01:00
|
|
|
}
|
|
|
|
|
if (m_socket.available() > 0) {
|
2017-05-27 20:27:48 +02:00
|
|
|
size_t received = m_socket.read_some(mBuffer.prepare(m_socket.available()));
|
2014-02-02 01:01:53 +01:00
|
|
|
if (received > 0) {
|
|
|
|
|
mBuffer.commit(received);
|
|
|
|
|
}
|
|
|
|
|
return received;
|
|
|
|
|
}
|
|
|
|
|
return 0;
|
|
|
|
|
}
|
|
|
|
|
|
2017-05-27 20:27:48 +02:00
|
|
|
size_t LocalStreamClientSocket::write()
|
2014-02-02 01:01:53 +01:00
|
|
|
{
|
|
|
|
|
if (!m_socket.is_open()) {
|
2017-05-27 20:27:48 +02:00
|
|
|
return 0;
|
2014-02-02 01:01:53 +01:00
|
|
|
}
|
|
|
|
|
if (mBuffer.size() == 0) {
|
|
|
|
|
return 0;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
auto size = boost::asio::write(m_socket, mBuffer.data(),
|
|
|
|
|
boost::asio::transfer_all());
|
|
|
|
|
mBuffer.consume(size);
|
|
|
|
|
return size;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
2017-07-01 19:36:11 +02:00
|
|
|
void AtlasStreamClient::output(const Element & item, size_t depth) const
|
2009-11-25 05:08:36 +00:00
|
|
|
{
|
2012-05-29 12:29:29 +01:00
|
|
|
output_element(std::cout, item, depth);
|
2009-11-25 05:08:36 +00:00
|
|
|
}
|
|
|
|
|
|
2012-06-09 15:51:58 -07:00
|
|
|
void AtlasStreamClient::output(const Root & ent) const
|
|
|
|
|
{
|
|
|
|
|
MapType entmap = ent->asMessage();
|
|
|
|
|
MapType::const_iterator Iend = entmap.end();
|
|
|
|
|
for (MapType::const_iterator I = entmap.begin(); I != Iend; ++I) {
|
|
|
|
|
const Element & item = I->second;
|
|
|
|
|
std::cout << std::string(spacing(), ' ') << I->first << ": ";
|
|
|
|
|
output(item, 1);
|
|
|
|
|
std::cout << std::endl;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2009-11-27 14:39:42 +00:00
|
|
|
/// \brief Function call from the base class when an object arrives from the
|
|
|
|
|
/// server
|
|
|
|
|
///
|
|
|
|
|
/// @param obj Object that has arrived from the server
|
|
|
|
|
void AtlasStreamClient::objectArrived(const Root & obj)
|
|
|
|
|
{
|
|
|
|
|
RootOperation op = Atlas::Objects::smart_dynamic_cast<RootOperation>(obj);
|
|
|
|
|
if (!op.isValid()) {
|
|
|
|
|
std::cerr << "ERROR: Non op object received from server"
|
|
|
|
|
<< std::endl << std::flush;;
|
2017-07-24 13:39:22 +02:00
|
|
|
if (!obj->isDefaultParent()) {
|
2009-11-27 14:39:42 +00:00
|
|
|
std::cerr << "NOTICE: Unexpected object has parent "
|
2017-07-24 13:39:22 +02:00
|
|
|
<< obj->getParent()
|
2009-11-27 14:39:42 +00:00
|
|
|
<< std::endl << std::flush;
|
|
|
|
|
}
|
|
|
|
|
if (!obj->isDefaultObjtype()) {
|
|
|
|
|
std::cerr << "NOTICE: Unexpected object has objtype "
|
|
|
|
|
<< obj->getObjtype()
|
|
|
|
|
<< std::endl << std::flush;
|
|
|
|
|
}
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
|
2014-02-04 22:24:27 +01:00
|
|
|
mOps.push_back(op);
|
2009-11-27 14:39:42 +00:00
|
|
|
}
|
|
|
|
|
|
2014-02-04 22:24:27 +01:00
|
|
|
void AtlasStreamClient::dispatch()
|
|
|
|
|
{
|
|
|
|
|
for (auto& op : mOps) {
|
|
|
|
|
operation(op);
|
|
|
|
|
}
|
|
|
|
|
mOps.clear();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
2009-11-27 15:09:44 +00:00
|
|
|
void AtlasStreamClient::operation(const RootOperation & op)
|
|
|
|
|
{
|
2016-07-09 23:30:05 +02:00
|
|
|
if (debug_flag) {
|
|
|
|
|
debug_print("Received:");
|
|
|
|
|
debug_dump(op, std::cout);
|
|
|
|
|
}
|
2018-09-22 23:11:22 +02:00
|
|
|
if (m_currentTask != nullptr) {
|
2009-11-28 18:31:28 +00:00
|
|
|
OpVector res;
|
|
|
|
|
m_currentTask->operation(op, res);
|
|
|
|
|
OpVector::const_iterator Iend = res.end();
|
|
|
|
|
for (OpVector::const_iterator I = res.begin(); I != Iend; ++I) {
|
|
|
|
|
send(*I);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if (m_currentTask->isComplete()) {
|
2018-09-22 23:11:22 +02:00
|
|
|
m_currentTask = nullptr;
|
2009-11-28 18:31:28 +00:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2009-11-28 18:14:49 +00:00
|
|
|
switch (op->getClassNo()) {
|
|
|
|
|
case Atlas::Objects::Operation::APPEARANCE_NO:
|
|
|
|
|
appearanceArrived(op);
|
|
|
|
|
break;
|
|
|
|
|
case Atlas::Objects::Operation::DISAPPEARANCE_NO:
|
|
|
|
|
disappearanceArrived(op);
|
|
|
|
|
break;
|
|
|
|
|
case Atlas::Objects::Operation::INFO_NO:
|
|
|
|
|
infoArrived(op);
|
|
|
|
|
break;
|
|
|
|
|
case Atlas::Objects::Operation::ERROR_NO:
|
|
|
|
|
errorArrived(op);
|
|
|
|
|
break;
|
|
|
|
|
case Atlas::Objects::Operation::SIGHT_NO:
|
|
|
|
|
sightArrived(op);
|
|
|
|
|
break;
|
|
|
|
|
case Atlas::Objects::Operation::SOUND_NO:
|
|
|
|
|
soundArrived(op);
|
|
|
|
|
break;
|
|
|
|
|
default:
|
|
|
|
|
break;
|
2009-11-27 15:09:44 +00:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2009-11-28 15:11:44 +00:00
|
|
|
void AtlasStreamClient::infoArrived(const RootOperation & op)
|
|
|
|
|
{
|
|
|
|
|
reply_flag = true;
|
|
|
|
|
if (!op->isDefaultFrom()) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
if (op->isDefaultArgs() || op->getArgs().empty()) {
|
|
|
|
|
std::cerr << "WARNING: Malformed account from server" << std::endl << std::flush;
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
if (op->isDefaultRefno()) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
if (op->getRefno() != serialNo) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
m_infoReply = op->getArgs().front();
|
|
|
|
|
}
|
|
|
|
|
|
2009-11-28 18:14:49 +00:00
|
|
|
void AtlasStreamClient::appearanceArrived(const RootOperation & op)
|
|
|
|
|
{
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void AtlasStreamClient::disappearanceArrived(const RootOperation & op)
|
|
|
|
|
{
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void AtlasStreamClient::sightArrived(const RootOperation & op)
|
|
|
|
|
{
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void AtlasStreamClient::soundArrived(const RootOperation & op)
|
|
|
|
|
{
|
|
|
|
|
}
|
|
|
|
|
|
2012-05-30 09:49:15 +01:00
|
|
|
void AtlasStreamClient::loginSuccess(const Atlas::Objects::Root & arg)
|
|
|
|
|
{
|
|
|
|
|
}
|
|
|
|
|
|
2009-11-27 15:09:44 +00:00
|
|
|
/// \brief Called when an Error operation arrives
|
|
|
|
|
///
|
|
|
|
|
/// @param op Operation to be processed
|
|
|
|
|
void AtlasStreamClient::errorArrived(const RootOperation & op)
|
|
|
|
|
{
|
|
|
|
|
reply_flag = true;
|
|
|
|
|
error_flag = true;
|
|
|
|
|
const std::vector<Root> & args = op->getArgs();
|
|
|
|
|
if (args.empty()) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
const Root & arg = args.front();
|
|
|
|
|
Element message_attr;
|
|
|
|
|
if (arg->copyAttr("message", message_attr) == 0 && message_attr.isString()) {
|
|
|
|
|
m_errorMessage = message_attr.String();
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2018-09-22 23:11:22 +02:00
|
|
|
AtlasStreamClient::AtlasStreamClient(boost::asio::io_service& io_service) :
|
|
|
|
|
m_io_service(io_service),
|
|
|
|
|
reply_flag(false),
|
|
|
|
|
error_flag(false),
|
|
|
|
|
serialNo(512),
|
|
|
|
|
m_socket(nullptr),
|
|
|
|
|
m_currentTask(nullptr),
|
|
|
|
|
m_spacing(2)
|
2009-11-24 18:36:45 +00:00
|
|
|
{
|
|
|
|
|
}
|
|
|
|
|
|
2018-10-04 21:50:03 +02:00
|
|
|
AtlasStreamClient::~AtlasStreamClient()
|
|
|
|
|
{
|
|
|
|
|
//Send a Logout op to inform the server that we're shutting down cleanly.
|
|
|
|
|
send(Atlas::Objects::Operation::Logout());
|
|
|
|
|
}
|
2009-11-24 18:36:45 +00:00
|
|
|
|
|
|
|
|
void AtlasStreamClient::send(const RootOperation & op)
|
|
|
|
|
{
|
2018-09-22 23:11:22 +02:00
|
|
|
if (!m_socket) {
|
2009-11-26 13:40:07 +00:00
|
|
|
return;
|
|
|
|
|
}
|
2012-05-29 18:18:36 +01:00
|
|
|
|
2016-07-09 23:30:05 +02:00
|
|
|
if (debug_flag) {
|
|
|
|
|
debug_print("Sending:");
|
|
|
|
|
debug_dump(op, std::cout);
|
|
|
|
|
}
|
|
|
|
|
|
2009-11-28 16:07:07 +00:00
|
|
|
reply_flag = false;
|
|
|
|
|
error_flag = false;
|
2014-02-02 01:01:53 +01:00
|
|
|
m_socket->getEncoder().streamObjectsMessage(op);
|
|
|
|
|
m_socket->write();
|
2009-11-24 18:36:45 +00:00
|
|
|
}
|
|
|
|
|
|
2017-07-01 19:36:11 +02:00
|
|
|
int AtlasStreamClient::connect(const std::string & host, unsigned short port)
|
2009-11-24 18:36:45 +00:00
|
|
|
{
|
2018-09-22 23:11:22 +02:00
|
|
|
m_socket.reset();
|
2014-02-02 01:01:53 +01:00
|
|
|
try {
|
2014-02-04 22:24:27 +01:00
|
|
|
std::function<void()> dispatcher = [&]{this->dispatch();};
|
2018-09-22 23:11:22 +02:00
|
|
|
m_socket = std::make_unique<TcpStreamClientSocket>(m_io_service, dispatcher, ip::tcp::endpoint(boost::asio::ip::address::from_string(host), port));
|
2014-02-02 01:01:53 +01:00
|
|
|
} catch (const std::exception& e) {
|
2009-11-24 18:36:45 +00:00
|
|
|
return -1;
|
|
|
|
|
}
|
2014-02-02 01:01:53 +01:00
|
|
|
return m_socket->negotiate(*this);
|
2009-11-24 18:36:45 +00:00
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
int AtlasStreamClient::connectLocal(const std::string & filename)
|
|
|
|
|
{
|
2018-09-22 23:11:22 +02:00
|
|
|
m_socket.reset();
|
|
|
|
|
|
2014-02-02 01:01:53 +01:00
|
|
|
try {
|
2014-02-04 22:24:27 +01:00
|
|
|
std::function<void()> dispatcher = [&]{this->dispatch();};
|
2018-09-22 23:11:22 +02:00
|
|
|
m_socket = std::make_unique<LocalStreamClientSocket>(m_io_service, dispatcher, local::stream_protocol::endpoint(filename));
|
2014-02-02 01:01:53 +01:00
|
|
|
} catch (const std::exception& e) {
|
2009-11-24 18:36:45 +00:00
|
|
|
return -1;
|
|
|
|
|
}
|
2014-02-02 01:01:53 +01:00
|
|
|
return m_socket->negotiate(*this);
|
2009-11-24 18:36:45 +00:00
|
|
|
}
|
|
|
|
|
|
2012-07-27 09:29:51 +01:00
|
|
|
int AtlasStreamClient::cleanDisconnect()
|
|
|
|
|
{
|
2018-09-22 23:11:22 +02:00
|
|
|
m_socket.reset();
|
2014-02-02 01:01:53 +01:00
|
|
|
// // Shutting down our write side will cause the server to get a HUP once
|
|
|
|
|
// // it has consumed all we have left for it
|
|
|
|
|
// m_ios->shutdown(true);
|
|
|
|
|
// // The server will then close the socket once we have all the responses
|
|
|
|
|
// while (this->poll(20, 0) == 0);
|
2012-07-27 09:29:51 +01:00
|
|
|
return 0;
|
|
|
|
|
}
|
|
|
|
|
|
2009-11-24 18:36:45 +00:00
|
|
|
|
|
|
|
|
int AtlasStreamClient::login(const std::string & username,
|
|
|
|
|
const std::string & password)
|
|
|
|
|
{
|
|
|
|
|
m_username = username;
|
|
|
|
|
|
|
|
|
|
Login l;
|
|
|
|
|
Anonymous account;
|
|
|
|
|
|
|
|
|
|
account->setAttr("username", username);
|
|
|
|
|
account->setAttr("password", password);
|
|
|
|
|
|
|
|
|
|
l->setArgs1(account);
|
2009-11-28 15:11:44 +00:00
|
|
|
l->setSerialno(newSerialNo());
|
2009-11-24 18:36:45 +00:00
|
|
|
|
|
|
|
|
send(l);
|
|
|
|
|
|
2012-05-24 11:07:22 +01:00
|
|
|
return waitForLoginResponse();
|
2009-11-24 18:36:45 +00:00
|
|
|
}
|
|
|
|
|
|
2009-12-03 19:03:59 +00:00
|
|
|
int AtlasStreamClient::create(const std::string & type,
|
|
|
|
|
const std::string & username,
|
|
|
|
|
const std::string & password)
|
2009-11-24 18:36:45 +00:00
|
|
|
{
|
|
|
|
|
m_username = username;
|
|
|
|
|
|
|
|
|
|
Create c;
|
|
|
|
|
Anonymous account;
|
|
|
|
|
|
|
|
|
|
account->setAttr("username", username);
|
|
|
|
|
account->setAttr("password", password);
|
2017-07-24 13:39:22 +02:00
|
|
|
account->setParent(type);
|
2009-11-24 18:36:45 +00:00
|
|
|
|
|
|
|
|
c->setArgs1(account);
|
2009-11-28 15:11:44 +00:00
|
|
|
c->setSerialno(newSerialNo());
|
2009-11-24 18:36:45 +00:00
|
|
|
|
|
|
|
|
send(c);
|
|
|
|
|
|
2012-05-24 11:07:22 +01:00
|
|
|
return waitForLoginResponse();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
int AtlasStreamClient::waitForLoginResponse()
|
|
|
|
|
{
|
2014-10-01 22:42:19 +02:00
|
|
|
while (poll(boost::posix_time::seconds(10)) == 0) {
|
2014-02-02 01:01:53 +01:00
|
|
|
if (reply_flag && !error_flag) {
|
|
|
|
|
if (m_infoReply->isDefaultId()) {
|
|
|
|
|
std::cerr << "Malformed reply" << std::endl << std::flush;
|
|
|
|
|
} else {
|
|
|
|
|
m_accountId = m_infoReply->getId();
|
2017-07-24 13:39:22 +02:00
|
|
|
m_accountType = m_infoReply->getParent();
|
2014-02-02 01:01:53 +01:00
|
|
|
loginSuccess(m_infoReply);
|
|
|
|
|
return 0;
|
|
|
|
|
}
|
|
|
|
|
reply_flag = false;
|
|
|
|
|
}
|
2009-11-28 15:11:44 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return -1;
|
2009-11-24 18:36:45 +00:00
|
|
|
}
|
2009-11-26 20:16:45 +00:00
|
|
|
|
2015-05-05 22:01:11 +02:00
|
|
|
int AtlasStreamClient::pollOne(const boost::posix_time::time_duration& duration)
|
|
|
|
|
{
|
|
|
|
|
if (!m_socket) {
|
|
|
|
|
return -1;
|
|
|
|
|
}
|
|
|
|
|
int result = m_socket->poll(duration, [&]()->bool{return !mOps.empty();});
|
|
|
|
|
if (result == -1) {
|
|
|
|
|
std::cerr << "Server disconnected" << std::endl << std::flush;
|
|
|
|
|
}
|
|
|
|
|
return result;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
2014-10-01 22:42:19 +02:00
|
|
|
int AtlasStreamClient::poll(const boost::posix_time::time_duration& duration)
|
2009-11-26 20:16:45 +00:00
|
|
|
{
|
2014-02-04 10:57:28 +01:00
|
|
|
if (!m_socket) {
|
|
|
|
|
return -1;
|
|
|
|
|
}
|
2015-05-05 22:01:11 +02:00
|
|
|
int result = m_socket->poll(duration, [](){return false;});
|
2014-02-02 01:01:53 +01:00
|
|
|
if (result == -1) {
|
|
|
|
|
std::cerr << "Server disconnected" << std::endl << std::flush;
|
2009-11-26 20:16:45 +00:00
|
|
|
}
|
2014-02-02 01:01:53 +01:00
|
|
|
return result;
|
|
|
|
|
}
|
2009-11-26 20:16:45 +00:00
|
|
|
|
2014-10-01 22:48:22 +02:00
|
|
|
int AtlasStreamClient::poll(int seconds, int microseconds)
|
2014-02-02 01:01:53 +01:00
|
|
|
{
|
2014-10-01 22:48:22 +02:00
|
|
|
return poll(boost::posix_time::seconds(seconds) + boost::posix_time::microseconds(microseconds));
|
2009-11-26 20:16:45 +00:00
|
|
|
}
|
2009-11-28 18:31:28 +00:00
|
|
|
|
2018-10-04 17:58:17 +02:00
|
|
|
int AtlasStreamClient::runTask(std::shared_ptr<ClientTask> task, const std::string & arg)
|
2009-11-28 18:31:28 +00:00
|
|
|
{
|
2018-09-22 23:11:22 +02:00
|
|
|
assert(task != nullptr);
|
2009-11-28 18:31:28 +00:00
|
|
|
|
2018-09-22 23:11:22 +02:00
|
|
|
if (m_currentTask != nullptr) {
|
2009-11-28 18:31:28 +00:00
|
|
|
std::cout << "Busy" << std::endl << std::flush;
|
|
|
|
|
return -1;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
m_currentTask = task;
|
|
|
|
|
|
|
|
|
|
OpVector res;
|
|
|
|
|
|
|
|
|
|
m_currentTask->setup(arg, res);
|
|
|
|
|
|
2009-12-21 22:06:49 +00:00
|
|
|
if (m_currentTask->isComplete()) {
|
2018-09-22 23:11:22 +02:00
|
|
|
m_currentTask = nullptr;
|
2009-12-21 22:06:49 +00:00
|
|
|
return -1;
|
|
|
|
|
}
|
|
|
|
|
|
2009-11-28 18:31:28 +00:00
|
|
|
OpVector::const_iterator Iend = res.end();
|
|
|
|
|
for (OpVector::const_iterator I = res.begin(); I != Iend; ++I) {
|
|
|
|
|
send(*I);
|
|
|
|
|
}
|
|
|
|
|
return 0;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
int AtlasStreamClient::endTask()
|
|
|
|
|
{
|
2018-09-22 23:11:22 +02:00
|
|
|
if (m_currentTask == nullptr) {
|
2009-11-28 18:31:28 +00:00
|
|
|
return -1;
|
|
|
|
|
}
|
2018-09-22 23:11:22 +02:00
|
|
|
m_currentTask = nullptr;
|
2009-11-28 18:31:28 +00:00
|
|
|
return 0;
|
|
|
|
|
}
|
2013-11-13 16:13:11 +01:00
|
|
|
|
|
|
|
|
bool AtlasStreamClient::hasTask() const {
|
|
|
|
|
return m_currentTask != nullptr;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
int AtlasStreamClient::pollUntilTaskComplete()
|
|
|
|
|
{
|
|
|
|
|
if (m_currentTask == nullptr) {
|
|
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
while (m_currentTask != nullptr) {
|
2014-10-01 22:42:19 +02:00
|
|
|
if (poll(boost::posix_time::milliseconds(100)) == -1) {
|
2013-11-13 16:13:11 +01:00
|
|
|
return -1;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
return 0;
|
|
|
|
|
}
|
|
|
|
|
|
2018-09-22 23:11:22 +02:00
|
|
|
|
2013-11-13 16:13:11 +01:00
|
|
|
|