aqhome: Prepared reorganizing IPC and nodes code around built-in event2 api.

This commit is contained in:
Martin Preuss
2025-02-26 00:49:33 +01:00
parent cf8edbbd5f
commit f63079af11
54 changed files with 2390 additions and 202 deletions

View File

@@ -48,7 +48,6 @@ static void GWENHYWFAR_CB _freeData(void *bp, void *p);
static int _handleSignal(AQH_OBJECT *o, uint32_t slotId, AQH_OBJECT *senderObject, int param1, void *param2);
static int _handleSocketReady(AQH_OBJECT *o);
static int _createListeningSocket(const char *addr, int port);
static int _socketSetReuseAddress(int sk, int fl);
static int _socketSetBlocking(int sk, int fl);
static int _translateHError(int herr);
@@ -63,7 +62,7 @@ static int _acceptConnection(int serverSocket);
* ------------------------------------------------------------------------------------------------
*/
AQH_OBJECT *AQH_TcpdObject_new(AQH_EVENT_LOOP *eventLoop)
AQH_OBJECT *AQH_TcpdObject_new(AQH_EVENT_LOOP *eventLoop, AQH_OBJECT *fdObject)
{
AQH_OBJECT *o;
AQH_TCPD_OBJECT *xo;
@@ -71,10 +70,18 @@ AQH_OBJECT *AQH_TcpdObject_new(AQH_EVENT_LOOP *eventLoop)
o=AQH_Object_new(eventLoop);
GWEN_NEW_OBJECT(AQH_TCPD_OBJECT, xo);
GWEN_INHERIT_SETDATA(AQH_OBJECT, AQH_TCPD_OBJECT, o, xo, _freeData);
xo->fdSocket=-1;
AQH_Object_SetSignalHandlerFn(o, _handleSignal);
xo->fdObject=fdObject;
xo->fdSocket=AQH_FdObject_GetFd(fdObject);
#if 0
/* create object for readable socket, connect to THIS, enable */
xo->fdObject=AQH_FdObject_new(AQH_Object_GetEventLoop(o), rv, AQH_FDOBJECT_FDMODE_READ);
AQH_Object_AddLink(xo->fdObject, AQH_FDOBJECT_SIGNAL_ISREADY, AQH_TCPD_OBJECT_SLOT_SOCKETREADY, o);
AQH_Object_Enable(xo->fdObject);
#endif
return o;
}
@@ -85,111 +92,13 @@ void GWENHYWFAR_CB _freeData(GWEN_UNUSED void *bp, void *p)
AQH_TCPD_OBJECT *xo;
xo=(AQH_TCPD_OBJECT*) p;
if (xo->fdObject)
AQH_Object_free(xo->fdObject);
if (xo->fdSocket)
close(xo->fdSocket);
GWEN_FREE_OBJECT(xo);
}
int AQH_TcpdObject_StartListening(AQH_OBJECT *o, const char *addr, int port)
{
if (o) {
AQH_TCPD_OBJECT *xo;
xo=GWEN_INHERIT_GETDATA(AQH_OBJECT, AQH_TCPD_OBJECT, o);
if (xo) {
if (xo->fdSocket<0) {
int rv;
rv=_createListeningSocket(addr, port);
if (rv<0) {
DBG_INFO(AQH_LOGDOMAIN, "here (%d)", rv);
return rv;
}
xo->fdSocket=rv;
/* create object for readable socket, connect to THIS, enable */
xo->fdObject=AQH_FdObject_new(AQH_Object_GetEventLoop(o), rv, AQH_FDOBJECT_FDMODE_READ);
AQH_Object_AddLink(xo->fdObject, AQH_FDOBJECT_SIGNAL_ISREADY, AQH_TCPD_OBJECT_SLOT_SOCKETREADY, o);
AQH_Object_Enable(xo->fdObject);
}
}
}
return GWEN_ERROR_GENERIC;
}
void AQH_TcpdObject_StopListening(AQH_OBJECT *o)
{
if (o) {
AQH_TCPD_OBJECT *xo;
xo=GWEN_INHERIT_GETDATA(AQH_OBJECT, AQH_TCPD_OBJECT, o);
if (xo) {
if (xo->fdObject) {
AQH_Object_Disable(xo->fdObject);
AQH_Object_free(xo->fdObject);
xo->fdObject=NULL;
}
if (xo->fdSocket<0) {
close(xo->fdSocket);
xo->fdSocket=-1;
}
}
}
}
int _handleSignal(AQH_OBJECT *o,
uint32_t slotId,
GWEN_UNUSED AQH_OBJECT *senderObject,
GWEN_UNUSED int param1,
GWEN_UNUSED void *param2)
{
switch(slotId) {
case AQH_TCPD_OBJECT_SLOT_SOCKETREADY: return _handleSocketReady(o);
default:
break;
}
return 0; /* not handled */
}
int _handleSocketReady(AQH_OBJECT *o)
{
AQH_TCPD_OBJECT *xo;
xo=GWEN_INHERIT_GETDATA(AQH_OBJECT, AQH_TCPD_OBJECT, o);
if (xo) {
int clientSk;
clientSk=_acceptConnection(xo->fdSocket);
if (clientSk<0) {
DBG_INFO(AQH_LOGDOMAIN, "here (%d)", clientSk);
}
else {
if (0==AQH_Object_EmitSignal(o, AQH_TCPD_OBJECT_SIGNAL_NEWCONN, clientSk, NULL)) {
DBG_INFO(AQH_LOGDOMAIN, "New connection not handled");
close(clientSk);
}
}
return 1; /* handled */
}
return 0; /* not handled */
}
int _createListeningSocket(const char *addr, int port)
int AQH_TcpdObject_CreateListeningSocket(const char *addr, int port)
{
int sk;
struct sockaddr_in inetAddr;
@@ -245,6 +154,48 @@ int _createListeningSocket(const char *addr, int port)
int _handleSignal(AQH_OBJECT *o,
uint32_t slotId,
GWEN_UNUSED AQH_OBJECT *senderObject,
GWEN_UNUSED int param1,
GWEN_UNUSED void *param2)
{
switch(slotId) {
case AQH_TCPD_OBJECT_SLOT_SOCKETREADY: return _handleSocketReady(o);
default:
break;
}
return 0; /* not handled */
}
int _handleSocketReady(AQH_OBJECT *o)
{
AQH_TCPD_OBJECT *xo;
xo=GWEN_INHERIT_GETDATA(AQH_OBJECT, AQH_TCPD_OBJECT, o);
if (xo) {
int clientSk;
clientSk=_acceptConnection(xo->fdSocket);
if (clientSk<0) {
DBG_INFO(AQH_LOGDOMAIN, "here (%d)", clientSk);
}
else {
if (0==AQH_Object_EmitSignal(o, AQH_TCPD_OBJECT_SIGNAL_NEWCONN, clientSk, NULL)) {
DBG_INFO(AQH_LOGDOMAIN, "New connection not handled");
close(clientSk);
}
}
return 1; /* handled */
}
return 0; /* not handled */
}
int _setHostAddr(struct in_addr *inetAddr, const char *sAddr)
{
inetAddr->s_addr=0;