Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion src/apps/towercalculator/KeyboardReader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@

#ifndef DOXYGEN_SHOULD_SKIP_THIS

#include "log/SemanticLogger.h"
#include "utils/Timeval.h"

#include <iostream>
Expand All @@ -56,7 +57,11 @@
namespace apps::towercalculator {

KeyboardReader::KeyboardReader(const std::function<void(long)>& cb)
: core::eventreceiver::ReadEventReceiver("KeyboardReader", 0)
: core::eventreceiver::ReadEventReceiver(
"KeyboardReader",
logger::LogScope{
logger::LogOrigin::Application, logger::LogBoundary::Application, "app", "towercalculator", logger::LogRole::Unknown, {}},
0)
, callBack(cb) {
if (!enable(STDIN_FILENO)) {
std::cout << "KeyboardReader not activated";
Expand Down
30 changes: 21 additions & 9 deletions src/core/DescriptorEventReceiver.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@

#include <cerrno>
#include <climits>
#include <optional>
#include <utility>

#endif /* DOXYGEN_SHOULD_SKIP_THIS */
Expand Down Expand Up @@ -84,6 +85,17 @@ namespace core {
, initialTimeout(timeout) {
}

DescriptorEventReceiver::DescriptorEventReceiver(const std::string& name,
DescriptorEventPublisher& descriptorEventPublisher,
logger::LogScope logScope,
const utils::Timeval& timeout)
: EventReceiver(name)
, descriptorEventPublisher(descriptorEventPublisher)
, logScope(logger::LogScopeOwner::fromScope(logScope))
, maxInactivity(timeout)
, initialTimeout(timeout) {
}

int DescriptorEventReceiver::getRegisteredFd() const {
return observedFd;
}
Expand All @@ -107,20 +119,20 @@ namespace core {

bool DescriptorEventReceiver::enable(int fd) {
if (enabled) {
log().warn("{}: Double enable", getName());
log().warn("{} descriptor: Double enable", descriptorEventPublisher.getName());
return false;
}

observedFd = fd;
if (descriptorEventPublisher.enable(this)) {
enabled = true;
log().trace("{}: Enabled", getName());
log().trace("{} descriptor enabled", descriptorEventPublisher.getName());
return true;
}

const int registrationError = errno != 0 ? errno : EIO;
observedFd = -1;
log().error("{}: Descriptor registration failed: fd={}", getName(), fd);
log().error("{} descriptor registration failed: fd={}", descriptorEventPublisher.getName(), fd);
errno = registrationError;
return false;
}
Expand All @@ -129,9 +141,9 @@ namespace core {
if (enabled) {
enabled = false;
descriptorEventPublisher.disable(this);
log().trace("{}: Disabled", getName());
log().trace("{} descriptor disabled", descriptorEventPublisher.getName());
} else {
log().warn("{}: Double disable", getName());
log().warn("{} descriptor: Double disable", descriptorEventPublisher.getName());
}
}

Expand All @@ -141,10 +153,10 @@ namespace core {
suspended = true;
descriptorEventPublisher.suspend(this);
} else {
log().warn("{}: Double suspend", getName());
log().warn("{} descriptor: Double suspend", descriptorEventPublisher.getName());
}
} else {
log().warn("{}: Suspend while not enabled", getName());
log().warn("{} descriptor: Suspend while not enabled", descriptorEventPublisher.getName());
}
}

Expand All @@ -155,10 +167,10 @@ namespace core {
lastTriggered = utils::Timeval::currentTime();
descriptorEventPublisher.resume(this);
} else {
log().warn("{}: Double resume", getName());
log().warn("{} descriptor: Double resume", descriptorEventPublisher.getName());
}
} else {
log().warn("{}: Resume while not enabled", getName());
log().warn("{} descriptor: Resume while not enabled", descriptorEventPublisher.getName());
}
}

Expand Down
5 changes: 5 additions & 0 deletions src/core/DescriptorEventReceiver.h
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@
#include "core/EventReceiver.h" // IWYU pragma: export
#include "core/Shutdown.h" // IWYU pragma: export
#include "log/LogScopeOwner.h"
#include "log/SemanticLogger.h"

namespace core {
class DescriptorEventPublisher;
Expand Down Expand Up @@ -92,6 +93,10 @@ namespace core {
DescriptorEventReceiver(const std::string& name,
DescriptorEventPublisher& descriptorEventPublisher,
const utils::Timeval& timeout = TIMEOUT::DISABLE);
DescriptorEventReceiver(const std::string& name,
DescriptorEventPublisher& descriptorEventPublisher,
logger::LogScope logScope,
const utils::Timeval& timeout = TIMEOUT::DISABLE);

public:
int getRegisteredFd() const;
Expand Down
6 changes: 5 additions & 1 deletion src/core/eventreceiver/ExceptionalConditionEventReceiver.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -43,17 +43,21 @@

#include "core/EventLoop.h"
#include "core/EventMultiplexer.h"
#include "log/SemanticLogger.h"

#ifndef DOXYGEN_SHOULD_SKIP_THIS

#endif /* DOXYGEN_SHOULD_SKIP_THIS */

namespace core::eventreceiver {

ExceptionalConditionEventReceiver::ExceptionalConditionEventReceiver(const std::string& name, const utils::Timeval& timeout)
ExceptionalConditionEventReceiver::ExceptionalConditionEventReceiver(const std::string& name,
logger::LogScope logScope,
const utils::Timeval& timeout)
: core::DescriptorEventReceiver(
name + " out of band",
core::EventLoop::instance().getEventMultiplexer().getDescriptorEventPublisher(core::EventMultiplexer::DISP_TYPE::EX),
logScope,
timeout) {
}

Expand Down
8 changes: 7 additions & 1 deletion src/core/eventreceiver/ExceptionalConditionEventReceiver.h
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,10 @@

#include "core/DescriptorEventReceiver.h" // IWYU pragma: export

namespace logger {
struct LogScope;
}

#ifndef DOXYGEN_SHOULD_SKIP_THIS

#include "utils/Timeval.h"
Expand All @@ -58,7 +62,9 @@ namespace core::eventreceiver {

class ExceptionalConditionEventReceiver : public core::DescriptorEventReceiver {
protected:
ExceptionalConditionEventReceiver(const std::string& name, const utils::Timeval& timeout = MAX_OUTOFBAND_INACTIVITY);
ExceptionalConditionEventReceiver(const std::string& name,
logger::LogScope logScope,
const utils::Timeval& timeout = MAX_OUTOFBAND_INACTIVITY);

virtual void outOfBandTimeout();

Expand Down
4 changes: 3 additions & 1 deletion src/core/eventreceiver/ReadEventReceiver.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -43,17 +43,19 @@

#include "core/EventLoop.h"
#include "core/EventMultiplexer.h"
#include "log/SemanticLogger.h"

#ifndef DOXYGEN_SHOULD_SKIP_THIS

#endif /* DOXYGEN_SHOULD_SKIP_THIS */

namespace core::eventreceiver {

ReadEventReceiver::ReadEventReceiver(const std::string& name, const utils::Timeval& timeout)
ReadEventReceiver::ReadEventReceiver(const std::string& name, logger::LogScope logScope, const utils::Timeval& timeout)
: core::DescriptorEventReceiver(
name + " read",
core::EventLoop::instance().getEventMultiplexer().getDescriptorEventPublisher(core::EventMultiplexer::DISP_TYPE::RD),
logScope,
timeout) {
}

Expand Down
6 changes: 5 additions & 1 deletion src/core/eventreceiver/ReadEventReceiver.h
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,10 @@

#include "core/DescriptorEventReceiver.h" // IWYU pragma: export

namespace logger {
struct LogScope;
}

#ifndef DOXYGEN_SHOULD_SKIP_THIS

#include "utils/Timeval.h"
Expand All @@ -56,7 +60,7 @@ namespace core::eventreceiver {

class ReadEventReceiver : public core::DescriptorEventReceiver {
protected:
ReadEventReceiver(const std::string& name, const utils::Timeval& timeout);
ReadEventReceiver(const std::string& name, logger::LogScope logScope, const utils::Timeval& timeout);

virtual void readTimeout();

Expand Down
4 changes: 3 additions & 1 deletion src/core/eventreceiver/WriteEventReceiver.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -43,17 +43,19 @@

#include "core/EventLoop.h"
#include "core/EventMultiplexer.h"
#include "log/SemanticLogger.h"

#ifndef DOXYGEN_SHOULD_SKIP_THIS

#endif /* DOXYGEN_SHOULD_SKIP_THIS */

namespace core::eventreceiver {

WriteEventReceiver::WriteEventReceiver(const std::string& name, const utils::Timeval& timeout)
WriteEventReceiver::WriteEventReceiver(const std::string& name, logger::LogScope logScope, const utils::Timeval& timeout)
: core::DescriptorEventReceiver(
name + " write",
core::EventLoop::instance().getEventMultiplexer().getDescriptorEventPublisher(core::EventMultiplexer::DISP_TYPE::WR),
logScope,
timeout) {
}

Expand Down
6 changes: 5 additions & 1 deletion src/core/eventreceiver/WriteEventReceiver.h
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,10 @@

#include "core/DescriptorEventReceiver.h" // IWYU pragma: export

namespace logger {
struct LogScope;
}

#ifndef DOXYGEN_SHOULD_SKIP_THIS

#include "utils/Timeval.h"
Expand All @@ -56,7 +60,7 @@ namespace core::eventreceiver {

class WriteEventReceiver : public core::DescriptorEventReceiver {
protected:
WriteEventReceiver(const std::string& name, const utils::Timeval& timeout);
WriteEventReceiver(const std::string& name, logger::LogScope logScope, const utils::Timeval& timeout);

virtual void writeTimeout();

Expand Down
44 changes: 36 additions & 8 deletions src/core/pipe/Pipe.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -47,9 +47,13 @@
#ifndef DOXYGEN_SHOULD_SKIP_THIS

#include "core/system/unistd.h"
#include "log/LogScopeOwner.h"
#include "log/SemanticLogger.h"
#include "utils/Timeval.h"

#include <atomic>
#include <cerrno>
#include <optional>
#include <system_error>
#include <utility>

Expand All @@ -58,6 +62,11 @@
namespace core::pipe {

namespace {
std::uint64_t allocateConnectionId() noexcept {
static std::atomic<std::uint64_t> nextConnectionId{1};
return nextConnectionId.fetch_add(1, std::memory_order_relaxed);
}

void closeDescriptor(int& fd) noexcept {
const int descriptor = std::exchange(fd, -1);
if (descriptor >= 0) {
Expand Down Expand Up @@ -90,11 +99,13 @@ namespace core::pipe {
}
} // namespace

Pipe::Pipe() noexcept
: Pipe(O_CLOEXEC) {
Pipe::Pipe(const std::string& instanceName)
: Pipe(O_CLOEXEC, instanceName) {
}

Pipe::Pipe(int flags) noexcept {
Pipe::Pipe(int flags, const std::string& instanceName)
: instanceName(instanceName)
, connectionId(allocateConnectionId()) {
int descriptors[2] = {-1, -1};
if (core::system::pipe2(descriptors, flags) != 0) {
error = errno;
Expand All @@ -106,7 +117,9 @@ namespace core::pipe {
}

Pipe::Pipe(Pipe&& pipe) noexcept
: readFd(std::exchange(pipe.readFd, -1))
: instanceName(std::move(pipe.instanceName))
, connectionId(pipe.connectionId)
, readFd(std::exchange(pipe.readFd, -1))
, writeFd(std::exchange(pipe.writeFd, -1))
, error(std::exchange(pipe.error, 0)) {
}
Expand All @@ -120,15 +133,19 @@ namespace core::pipe {
if (this != &pipe) {
closeRead();
closeWrite();
instanceName = std::move(pipe.instanceName);
connectionId = pipe.connectionId;
readFd = std::exchange(pipe.readFd, -1);
writeFd = std::exchange(pipe.writeFd, -1);
error = std::exchange(pipe.error, 0);
}
return *this;
}

Pipe::Pipe(const std::function<void(PipeSource&, PipeSink&)>& onSuccess, const std::function<void(int)>& onError)
: Pipe(O_NONBLOCK | O_CLOEXEC) {
Pipe::Pipe(const std::function<void(PipeSource&, PipeSink&)>& onSuccess,
const std::function<void(int)>& onError,
const std::string& instanceName)
: Pipe(O_NONBLOCK | O_CLOEXEC, instanceName) {
if (!hasReadFd() || !hasWriteFd()) {
onError(error);
return;
Expand Down Expand Up @@ -197,6 +214,15 @@ namespace core::pipe {
closeDescriptor(writeFd);
}

logger::LogScopeOwner Pipe::makeLogScope() const {
return logger::LogScopeOwner(logger::LogOrigin::Framework,
logger::LogBoundary::Connection,
"core.pipe",
instanceName.empty() ? std::nullopt : std::optional<std::string>(instanceName),
std::nullopt,
std::to_string(connectionId));
}

PipeSink* Pipe::releaseReadAsSink() {
return releaseReadAsSink(PipeSink::DEFAULT_MAX_BYTES_PER_EVENT, utils::Timeval({60, 0}));
}
Expand All @@ -209,7 +235,8 @@ namespace core::pipe {
makeNonBlocking(readFd);
const int descriptor = releaseReadFd();
try {
return new PipeSink(descriptor, maxBytesPerEvent, timeout);
const logger::LogScopeOwner logScope = makeLogScope();
return new PipeSink(descriptor, logScope.scope(), maxBytesPerEvent, timeout);
} catch (...) {
readFd = descriptor;
throw;
Expand All @@ -228,7 +255,8 @@ namespace core::pipe {
makeNonBlocking(writeFd);
const int descriptor = releaseWriteFd();
try {
return new PipeSource(descriptor, maxQueuedBytes, timeout);
const logger::LogScopeOwner logScope = makeLogScope();
return new PipeSource(descriptor, logScope.scope(), maxQueuedBytes, timeout);
} catch (...) {
writeFd = descriptor;
throw;
Expand Down
Loading
Loading