stream: add a stream class abstracting BSD sockets
Currently only synchronous operation is supported, but this will be extended with asynchronous methods using the new watcher.
This commit is contained in:
@@ -26,7 +26,7 @@ credentials/sets/callback_cred.c credentials/auth_cfg.c database/database.c \
|
|||||||
database/database_factory.c fetcher/fetcher.c fetcher/fetcher_manager.c eap/eap.c \
|
database/database_factory.c fetcher/fetcher.c fetcher/fetcher_manager.c eap/eap.c \
|
||||||
ipsec/ipsec_types.c \
|
ipsec/ipsec_types.c \
|
||||||
networking/host.c networking/host_resolver.c networking/packet.c \
|
networking/host.c networking/host_resolver.c networking/packet.c \
|
||||||
networking/tun_device.c \
|
networking/tun_device.c networking/streams/stream.c \
|
||||||
pen/pen.c plugins/plugin_loader.c plugins/plugin_feature.c processing/jobs/job.c \
|
pen/pen.c plugins/plugin_loader.c plugins/plugin_feature.c processing/jobs/job.c \
|
||||||
processing/jobs/callback_job.c processing/processor.c processing/scheduler.c \
|
processing/jobs/callback_job.c processing/processor.c processing/scheduler.c \
|
||||||
processing/watcher.c resolver/resolver_manager.c resolver/rr_set.c \
|
processing/watcher.c resolver/resolver_manager.c resolver/rr_set.c \
|
||||||
|
|||||||
@@ -24,7 +24,7 @@ credentials/sets/callback_cred.c credentials/auth_cfg.c database/database.c \
|
|||||||
database/database_factory.c fetcher/fetcher.c fetcher/fetcher_manager.c eap/eap.c \
|
database/database_factory.c fetcher/fetcher.c fetcher/fetcher_manager.c eap/eap.c \
|
||||||
ipsec/ipsec_types.c \
|
ipsec/ipsec_types.c \
|
||||||
networking/host.c networking/host_resolver.c networking/packet.c \
|
networking/host.c networking/host_resolver.c networking/packet.c \
|
||||||
networking/tun_device.c \
|
networking/tun_device.c networking/streams/stream.c \
|
||||||
pen/pen.c plugins/plugin_loader.c plugins/plugin_feature.c processing/jobs/job.c \
|
pen/pen.c plugins/plugin_loader.c plugins/plugin_feature.c processing/jobs/job.c \
|
||||||
processing/jobs/callback_job.c processing/processor.c processing/scheduler.c \
|
processing/jobs/callback_job.c processing/processor.c processing/scheduler.c \
|
||||||
processing/watcher.c resolver/resolver_manager.c resolver/rr_set.c \
|
processing/watcher.c resolver/resolver_manager.c resolver/rr_set.c \
|
||||||
@@ -65,7 +65,7 @@ credentials/auth_cfg.h credentials/credential_set.h credentials/cert_validator.h
|
|||||||
database/database.h database/database_factory.h fetcher/fetcher.h \
|
database/database.h database/database_factory.h fetcher/fetcher.h \
|
||||||
fetcher/fetcher_manager.h eap/eap.h pen/pen.h ipsec/ipsec_types.h \
|
fetcher/fetcher_manager.h eap/eap.h pen/pen.h ipsec/ipsec_types.h \
|
||||||
networking/host.h networking/host_resolver.h networking/packet.h \
|
networking/host.h networking/host_resolver.h networking/packet.h \
|
||||||
networking/tun_device.h \
|
networking/tun_device.h networking/streams/stream.h \
|
||||||
resolver/resolver.h resolver/resolver_response.h resolver/rr_set.h \
|
resolver/resolver.h resolver/resolver_response.h resolver/rr_set.h \
|
||||||
resolver/rr.h resolver/resolver_manager.h \
|
resolver/rr.h resolver/resolver_manager.h \
|
||||||
plugins/plugin_loader.h plugins/plugin.h plugins/plugin_feature.h \
|
plugins/plugin_loader.h plugins/plugin.h plugins/plugin_feature.h \
|
||||||
|
|||||||
@@ -0,0 +1,119 @@
|
|||||||
|
/*
|
||||||
|
* Copyright (C) 2013 Martin Willi
|
||||||
|
* Copyright (C) 2013 revosec AG
|
||||||
|
*
|
||||||
|
* 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. See <http://www.fsf.org/copyleft/gpl.txt>.
|
||||||
|
*
|
||||||
|
* 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.
|
||||||
|
*/
|
||||||
|
|
||||||
|
#include "stream.h"
|
||||||
|
|
||||||
|
#include <errno.h>
|
||||||
|
#include <unistd.h>
|
||||||
|
|
||||||
|
typedef struct private_stream_t private_stream_t;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Private data of an stream_t object.
|
||||||
|
*/
|
||||||
|
struct private_stream_t {
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Public stream_t interface.
|
||||||
|
*/
|
||||||
|
stream_t public;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Underlying socket
|
||||||
|
*/
|
||||||
|
int fd;
|
||||||
|
};
|
||||||
|
|
||||||
|
METHOD(stream_t, read_, ssize_t,
|
||||||
|
private_stream_t *this, void *buf, size_t len, bool block)
|
||||||
|
{
|
||||||
|
while (TRUE)
|
||||||
|
{
|
||||||
|
ssize_t ret;
|
||||||
|
|
||||||
|
if (block)
|
||||||
|
{
|
||||||
|
ret = read(this->fd, buf, len);
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
ret = recv(this->fd, buf, len, MSG_DONTWAIT);
|
||||||
|
if (ret == -1 && errno == EAGAIN)
|
||||||
|
{
|
||||||
|
/* unify EGAIN and EWOULDBLOCK */
|
||||||
|
errno = EWOULDBLOCK;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (ret == -1 && errno == EINTR)
|
||||||
|
{ /* interrupted, try again */
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
return ret;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
METHOD(stream_t, write_, ssize_t,
|
||||||
|
private_stream_t *this, void *buf, size_t len, bool block)
|
||||||
|
{
|
||||||
|
ssize_t ret;
|
||||||
|
|
||||||
|
while (TRUE)
|
||||||
|
{
|
||||||
|
if (block)
|
||||||
|
{
|
||||||
|
ret = write(this->fd, buf, len);
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
ret = send(this->fd, buf, len, MSG_DONTWAIT);
|
||||||
|
if (ret == -1 && errno == EAGAIN)
|
||||||
|
{
|
||||||
|
/* unify EGAIN and EWOULDBLOCK */
|
||||||
|
errno = EWOULDBLOCK;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (ret == -1 && errno == EINTR)
|
||||||
|
{ /* interrupted, try again */
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
return ret;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
METHOD(stream_t, destroy, void,
|
||||||
|
private_stream_t *this)
|
||||||
|
{
|
||||||
|
close(this->fd);
|
||||||
|
free(this);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* See header
|
||||||
|
*/
|
||||||
|
stream_t *stream_create_from_fd(int fd)
|
||||||
|
{
|
||||||
|
private_stream_t *this;
|
||||||
|
|
||||||
|
INIT(this,
|
||||||
|
.public = {
|
||||||
|
.read = _read_,
|
||||||
|
.write = _write_,
|
||||||
|
.destroy = _destroy,
|
||||||
|
},
|
||||||
|
.fd = fd,
|
||||||
|
);
|
||||||
|
|
||||||
|
return &this->public;
|
||||||
|
}
|
||||||
@@ -0,0 +1,83 @@
|
|||||||
|
/*
|
||||||
|
* Copyright (C) 2013 Martin Willi
|
||||||
|
* Copyright (C) 2013 revosec AG
|
||||||
|
*
|
||||||
|
* 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. See <http://www.fsf.org/copyleft/gpl.txt>.
|
||||||
|
*
|
||||||
|
* 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.
|
||||||
|
*/
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @defgroup stream stream
|
||||||
|
* @{ @ingroup streams
|
||||||
|
*/
|
||||||
|
|
||||||
|
#ifndef STREAM_H_
|
||||||
|
#define STREAM_H_
|
||||||
|
|
||||||
|
typedef struct stream_t stream_t;
|
||||||
|
|
||||||
|
#include <library.h>
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Constructor function prototype for stream_t.
|
||||||
|
*
|
||||||
|
* @param uri URI to create a stream for
|
||||||
|
* @return stream instance, NULL on error
|
||||||
|
*/
|
||||||
|
typedef stream_t*(*stream_constructor_t)(char *uri);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Abstraction of a Berkley socket using stream semantics.
|
||||||
|
*/
|
||||||
|
struct stream_t {
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Read data from the stream.
|
||||||
|
*
|
||||||
|
* If "block" is FALSE and no data is available, the function returns -1
|
||||||
|
* and sets errno to EWOULDBLOCK.
|
||||||
|
*
|
||||||
|
* @param buf data buffer to read into
|
||||||
|
* @param len number of bytes to read
|
||||||
|
* @param block TRUE to use a blocking read
|
||||||
|
* @return number of bytes read, -1 on error
|
||||||
|
*/
|
||||||
|
ssize_t (*read)(stream_t *this, void *buf, size_t len, bool block);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Write data to the stream.
|
||||||
|
*
|
||||||
|
* If "block" is FALSE and the write would block, the function returns -1
|
||||||
|
* and sets errno to EWOULDBLOCK.
|
||||||
|
*
|
||||||
|
* @param buf data buffer to write
|
||||||
|
* @param len number of bytes to write
|
||||||
|
* @param block TRUE to use a blocking write
|
||||||
|
* @return number of bytes written, -1 on error
|
||||||
|
*/
|
||||||
|
ssize_t (*write)(stream_t *this, void *buf, size_t len, bool block);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Destroy a stream_t.
|
||||||
|
*/
|
||||||
|
void (*destroy)(stream_t *this);
|
||||||
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Create a stream from a file descriptor.
|
||||||
|
*
|
||||||
|
* The file descriptor MUST be a socket for non-blocking operation.
|
||||||
|
*
|
||||||
|
* @param fd file descriptor to wrap into a stream_t
|
||||||
|
* @return stream instance
|
||||||
|
*/
|
||||||
|
stream_t *stream_create_from_fd(int fd);
|
||||||
|
|
||||||
|
#endif /** STREAM_H_ @}*/
|
||||||
Reference in New Issue
Block a user