Skip to content
Draft
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
74 changes: 49 additions & 25 deletions core/src/main/java/io/jstach/rainbowgum/LogAppender.java
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
import java.util.List;
import java.util.Locale;
import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.locks.ReentrantLock;
Expand Down Expand Up @@ -410,9 +411,9 @@ private static LogAppender single(List<? extends LogAppender> appenders) {

}

interface AppenderVisitor {
interface AppenderVisitor<R> {

boolean consume(DirectLogAppender appender);
Optional<R> apply(DirectLogAppender appender);

}

Expand Down Expand Up @@ -524,10 +525,21 @@ static InternalLogAppender of(LogAppender appender) {

/**
* THIS IS A JAVADOC BUG.
* @param <R> return type.
* @param visitor ignore
* @return true if stop.
*/
boolean visit(AppenderVisitor visitor);
<R> Optional<R> visit(AppenderVisitor<R> visitor);

default <R> Optional<R> useOutput(String name, Function<? super LogOutput, ? extends R> action) {
return visit(da -> {
if (name.equals(da.name())) {
return Optional.ofNullable(action.apply(da.output()));
}
return Optional.empty();

});
}

// InternalLogAppender changeLock(AppenderLock lock);
//
Expand All @@ -543,6 +555,7 @@ static InternalLogAppender of(LogAppender appender) {

}

// This is super private because there is no locking protection
sealed interface DirectLogAppender extends InternalLogAppender {

String name();
Expand All @@ -551,6 +564,7 @@ sealed interface DirectLogAppender extends InternalLogAppender {

LogEncoder encoder();

// Should be called within lock block.
default List<LogResponse> _request(LogAction action) {
List<LogResponse> r = switch (action) {
case LogAction.StandardAction a -> switch (a) {
Expand All @@ -563,17 +577,17 @@ default List<LogResponse> _request(LogAction action) {
return r;
}

default LogResponse reopen() {
private LogResponse reopen() {
var status = output().reopen();
return new Response(LogOutput.class, name(), status);
}

default LogResponse flush() {
private LogResponse flush() {
output().flush();
return new Response(LogOutput.class, name(), LogResponse.Status.StandardStatus.OK);
}

default LogResponse status() {
private LogResponse status() {
Status status;
try {
status = output().status();
Expand All @@ -584,16 +598,6 @@ default LogResponse status() {
return new Response(LogOutput.class, name(), status);
}

static List<DirectLogAppender> findAppenders(ServiceRegistry registry) {
List<DirectLogAppender> appenders = new ArrayList<>();
for (var a : registry.find(LogAppender.class)) {
if (a instanceof InternalLogAppender internal) {
internal.visit(appenders::add);
}
}
return appenders;
}

static DirectLogAppender of(String name, LogOutput output, LogEncoder encoder,
Set<LogAppender.AppenderFlag> flags) {
var lock = AppenderLock.of(flags);
Expand All @@ -603,6 +607,8 @@ static DirectLogAppender of(String name, LogOutput output, LogEncoder encoder,
return new DefaultLogAppender(name, output, encoder, flags, lock);
}

// <R> R execute(Function<? super LogOutput, ? extends R> action);

// @Override
DirectLogAppender withFlags(Set<LogAppender.AppenderFlag> flags);

Expand Down Expand Up @@ -667,10 +673,10 @@ public String toString() {
+ flags + "]";
}

@Override
public boolean visit(AppenderVisitor visitor) {
return visitor.consume(this);
}
// @Override
// public <R> Optional<R> visit(AppenderVisitor<R> visitor) {
// return visitor.apply(this);
// }

@Override
public String name() {
Expand Down Expand Up @@ -753,13 +759,20 @@ default void start(LogConfig config) {
}

@Override
default boolean visit(AppenderVisitor visitor) {
for (var appender : components()) {
if (appender.visit(visitor)) {
return true;
default <R> Optional<R> visit(AppenderVisitor<R> visitor) {
lock().lock();
try {
for (var appender : components()) {
var o = appender.visit(visitor);
if (o.isPresent()) {
return o;
}
}
return Optional.empty();
}
finally {
lock().unlock();
}
return false;
}

@Override
Expand Down Expand Up @@ -839,6 +852,17 @@ public List<LogResponse> act(LogAction action) {
}
}

@Override
public <R> Optional<R> visit(AppenderVisitor<R> visitor) {
lock.lock();
try {
return visitor.apply(this);
}
finally {
lock.unlock();
}
}

@Override
public void close() {
lock.lock();
Expand Down
41 changes: 33 additions & 8 deletions core/src/main/java/io/jstach/rainbowgum/LogOutputRegistry.java
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@
import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Function;
import java.util.stream.Stream;

import io.jstach.rainbowgum.LogOutput.OutputProvider;
import io.jstach.rainbowgum.output.FileOutput;
Expand All @@ -33,11 +35,26 @@ public sealed interface LogOutputRegistry extends OutputProvider permits Default
* @param scheme URI scheme to match for.
* @param provider provider for scheme.
*/
public void register(String scheme, OutputProvider provider);
void register(String scheme, OutputProvider provider);

/**
* Finds an output by name.
* @param name output name.
* Allows you to perform a safe operation on the {@link LogOutput} corresponding to
* the named appender by calling the function within the appender locking mechanism.
* The output will not get any write calls through the publisher/appender while the
* passed in function is being called.
* @param <R> return type.
* @param name the appender name.
* @param action function to perform the operation. <strong>The function not share the
* output outside.</strong> If the function returns <code>null</code> the returned
* optional will always be empty but the function may have been applied so it is
* recommend you do not do this.
* @return user desired return if an appender is found.
*/
<R> Optional<? extends R> useOutput(String name, Function<? super LogOutput, ? extends R> action);

/**
* Finds an output by name of the appender.
* @param name of appender that owns output (and not the configuration name).
* @return maybe an output.
*/
Optional<LogOutput> output(String name);
Expand All @@ -49,21 +66,21 @@ public sealed interface LogOutputRegistry extends OutputProvider permits Default
* @return the output status of reopened outputs or an empty list if no outputs were
* reopened.
*/
public List<LogResponse> reopen();
List<LogResponse> reopen();

/**
* Attempts to flush all outputs usually for log rotation. This call will block if it
* attempts to flush. If flush is already happening an empty list will be returned.
* @return the output status of reopened outputs or an empty list if no outputs were
* reopened.
*/
public List<LogResponse> flush();
List<LogResponse> flush();

/**
* Will retrieve the status of all outputs usually for health checking.
* @return list of status of outputs.
*/
public List<LogResponse> status();
List<LogResponse> status();

}

Expand Down Expand Up @@ -107,6 +124,15 @@ public List<LogResponse> status() {
return _request(LogAction.StandardAction.STATUS);
}

@Override
public <R> Optional<? extends R> useOutput(String name, Function<? super LogOutput, ? extends R> action) {
return findAppenders().map(a -> a.useOutput(name, action)).flatMap(o -> o.stream()).findFirst();
}

private Stream<InternalLogAppender> findAppenders() {
return serviceRegistry.find(LogAppender.class).stream().map(a -> InternalLogAppender.of(a));
}

private List<LogResponse> requestIO(LogAction action) {
if (reopenLock.tryLock()) {
try {
Expand All @@ -125,8 +151,7 @@ private List<LogResponse> _request(LogAction action) {
/*
* TODO check rainbowgum is actually running.
*/
return Actor.act(serviceRegistry.find(LogAppender.class).stream().map(a -> InternalLogAppender.of(a)).toList(),
action);
return Actor.act(findAppenders().toList(), action);
}

@Override
Expand Down