diff --git a/direct/src/distributed/cConnectionRepository.cxx b/direct/src/distributed/cConnectionRepository.cxx index 596db24ab7..7c6673d509 100644 --- a/direct/src/distributed/cConnectionRepository.cxx +++ b/direct/src/distributed/cConnectionRepository.cxx @@ -27,6 +27,7 @@ #ifdef HAVE_PYTHON #include "py_panda.h" #include "dcClass_ext.h" +extern struct Dtool_PyTypedObject Dtool_DCClass; #endif using std::endl; @@ -46,6 +47,7 @@ CConnectionRepository(bool has_owner_view, bool threaded_net) : _lock("CConnectionRepository::_lock"), #ifdef HAVE_PYTHON _python_repository(nullptr), + _python_ai_datagramiterator(nullptr), #endif #ifdef HAVE_OPENSSL _http_conn(nullptr), @@ -77,6 +79,24 @@ CConnectionRepository(bool has_owner_view, bool threaded_net) : } #endif _tcp_header_size = tcp_header_size; + +#ifdef HAVE_PYTHON + Dtool_TypeMap *tmap = Dtool_GetGlobalTypeMap(); + auto tmap_search = tmap->find("DatagramIterator"); + // Make sure we found what we were looking for. + nassertv(tmap_search != tmap->end()); + Dtool_PyTypedObject *Dtool_DatagramIterator = tmap_search->second; + // If this is a nullptr. We can't go further. + nassertv(Dtool_DatagramIterator != nullptr); + PyObject *PyDitterator = DTool_CreatePyInstance(&_di, *Dtool_DatagramIterator, false, false); + if(PyDitterator != nullptr) { + _python_ai_datagramiterator = Py_BuildValue("(O)", PyDitterator); + if (PyErr_Occurred()) { + PyErr_Print(); + } + } + Py_XDECREF(PyDitterator); +#endif } /** @@ -1013,3 +1033,160 @@ describe_message(std::ostream &out, const string &prefix, } } } + +#ifdef HAVE_PYTHON +#ifdef WANT_NATIVE_NET + +bool CConnectionRepository::network_based_reader_and_yielder(PyObject *PycallBackFunction,ClockObject &clock, float returnBy) { + ReMutexHolder holder(_lock); + distributed_cat.debug() << "Doing Network Based Read.\n"; + while(is_connected()) { + check_datagram_ai(PycallBackFunction); + if(is_connected()) + _bdc.Flush(); + float currentTime = clock.get_real_time(); + float dif_time = returnBy - currentTime; + if(dif_time <= 0.001) // to avoi over runs.. + break; + if(is_connected()) + _bdc.WaitForNetworkReadEvent(dif_time); + } + return false; +} + +bool CConnectionRepository::check_datagram_ai(PyObject *PycallBackFunction) { + ReMutexHolder holder(_lock); +#if defined(HAVE_THREADS) && !defined(SIMPLE_THREADS) + PyGILState_STATE gstate; + gstate = PyGILState_Ensure(); +#endif + // these could be static .. not + PyObject *doId2do = nullptr; + float startTime = 0; + float endTime = 0; + // this seems weird...here + _bdc.Flush(); + while (_bdc.GetMessage(_dg)) { + if (get_verbose()) + describe_message(nout, "RECV", _dg); + + if (_time_warning > 0) + startTime = ClockObject::get_global_clock()->get_real_time(); + + // Start breaking apart the datagram. + _di.assign(_dg); + unsigned char wc_cnt = _di.get_uint8(); + _msg_channels.clear(); + for (unsigned char lp1 = 0; lp1 < wc_cnt; lp1++) { + _msg_channels.push_back(_di.get_uint64()); + } + + _msg_sender = _di.get_uint64(); + _msg_type = _di.get_uint16(); + + if (_msg_type == STATESERVER_OBJECT_UPDATE_FIELD) { + if (doId2do == nullptr) { + // this is my attemp to take it out of the inner loop RHH + doId2do = PyObject_GetAttrString(_python_repository, "doId2do"); + } + + bool success = handle_update_field_ai(doId2do); + + if (!success) { + Py_XDECREF(doId2do); + if (_time_warning > 0) { + endTime = ClockObject::get_global_clock()->get_real_time(); + if ( _time_warning < (endTime - startTime)) { + nout << "msg " << _msg_type <<" from " << _msg_sender << " took "<< (endTime-startTime) << "secs to process\n"; + _dg.dump_hex(nout,2); + } + } +#if defined(HAVE_THREADS) && !defined(SIMPLE_THREADS) + PyGILState_Release(gstate); +#endif + return false; + } + } else { + distributed_cat.debug() << "Calling function callback in check_datagram_ai()!\n"; + nassertr(_python_ai_datagramiterator != nullptr, false); + + PyObject *result = PyEval_CallObject(PycallBackFunction, _python_ai_datagramiterator); + + Py_XDECREF(result); + if (PyErr_Occurred()) { + Py_XDECREF(doId2do); + if (_time_warning > 0) { + endTime = ClockObject::get_global_clock()->get_real_time(); + if ( _time_warning < (endTime - startTime)) { + nout << "msg " << _msg_type <<" from " << _msg_sender << " took "<< (endTime-startTime) << "secs to process\n"; + _dg.dump_hex(nout,2); + } + } +#if defined(HAVE_THREADS) && !defined(SIMPLE_THREADS) + PyGILState_Release(gstate); +#endif + return true; + } + } + + if (_time_warning > 0) { + endTime = ClockObject::get_global_clock()->get_real_time(); + if ( _time_warning < (endTime - startTime)) { + nout << "msg " << _msg_type <<" from " << _msg_sender << " took "<< (endTime-startTime) << "secs to process\n"; + _dg.dump_hex(nout,2); + } + } + + } + + Py_XDECREF(doId2do); + +#if defined(HAVE_THREADS) && !defined(SIMPLE_THREADS) + PyGILState_Release(gstate); +#endif + return false; +} + +#endif // #ifdef WANT_NATIVE_NET +#endif // #ifdef HAVE_PYTHON + +#ifdef HAVE_PYTHON +#ifdef WANT_NATIVE_NET + +bool CConnectionRepository::handle_update_field_ai(PyObject *doId2do) { + PStatTimer timer(_update_pcollector); + unsigned int do_id = _di.get_uint32(); + +#ifdef USE_PYTHON_2_2_OR_EARLIER + PyObject *doId = PyInt_FromLong(do_id); +#else + PyObject *doId = PyLong_FromUnsignedLong(do_id); +#endif + PyObject *distobj = PyDict_GetItem(doId2do, doId); + Py_XDECREF(doId); + + if (distobj != nullptr) { + PyObject *dclass_obj = PyObject_GetAttrString(distobj, "dclass"); + nassertr(dclass_obj != nullptr, false); + + PyObject *dclass_obj_this = PyObject_GetAttrString(dclass_obj, "this"); + Py_XDECREF(dclass_obj); + nassertr(dclass_obj_this != nullptr, false); + + DCClass *dclass = (DCClass *)PyLong_AsVoidPtr(dclass_obj_this); + Py_XDECREF(dclass_obj_this); + nassertr(dclass != nullptr, false); + + Py_XINCREF(distobj); + invoke_extension(dclass).receive_update(distobj, _di); + Py_XDECREF(distobj); + + if (PyErr_Occurred()) { + return false; + } + } + return true; +} + +#endif // #ifdef WANT_NATIVE_NET +#endif // #ifdef HAVE_PYTHON diff --git a/direct/src/distributed/cConnectionRepository.h b/direct/src/distributed/cConnectionRepository.h index 2ee1d9d0cd..267056bea6 100644 --- a/direct/src/distributed/cConnectionRepository.h +++ b/direct/src/distributed/cConnectionRepository.h @@ -116,6 +116,12 @@ PUBLISHED: #endif BLOCKING bool check_datagram(); +#ifdef HAVE_PYTHON +#ifdef WANT_NATIVE_NET + BLOCKING bool check_datagram_ai(PyObject *PycallBackFunction); + BLOCKING bool network_based_reader_and_yielder(PyObject *PycallBackFunction, ClockObject &clock, float returnBy); +#endif +#endif BLOCKING INLINE void get_datagram(Datagram &dg); BLOCKING INLINE void get_datagram_iterator(DatagramIterator &di); @@ -160,6 +166,11 @@ PUBLISHED: INLINE float get_time_warning() const; private: +#ifdef HAVE_PYTHON +#ifdef WANT_NATIVE_NET + bool handle_update_field_ai(PyObject *doId2do); +#endif +#endif bool do_check_datagram(); bool handle_update_field(); bool handle_update_field_owner(); @@ -172,6 +183,7 @@ private: #ifdef HAVE_PYTHON PyObject *_python_repository; + PyObject *_python_ai_datagramiterator; #endif #ifdef HAVE_OPENSSL