diff --git a/AGENTS.md b/AGENTS.md index de2827f..cbf3e33 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -198,6 +198,10 @@ - Keep constructors to dependency and configuration assignment. Put startup parsing, object-graph construction, defaults, validation across components, and any other non-trivial creation logic in explicit composition factories. +- Implement technical boundary decoration once over JDK functional interfaces + and adapt business-named ports through method references. Do not make domain + ports extend `Consumer` or `Function` merely to simplify a decorator, and do + not create one decorator class for every port shape. - Test the shared application through its public API with injected adapters and deterministic dependencies. Give adapter implementations focused contract tests and give composition factories a small wiring smoke test. A trivial diff --git a/JOURNAL.md b/JOURNAL.md index 44f7462..5bb248d 100644 --- a/JOURNAL.md +++ b/JOURNAL.md @@ -201,3 +201,12 @@ maintenance updates. Direct-push CI now runs only for `main`; pull requests keep their own unit and integration run without a duplicate branch-push run. The license file matches GitHub's canonical Apache 2.0 template exactly, and hosting metadata identifies it as `Apache-2.0`. + +## 2026-08-10 — Consolidate functional boundary observation + +Replaced one observability class per port with a single package-private +operation observer over JDK `Consumer`, `IntFunction`, `Supplier`, and +`ToLongFunction`. The public observability factory adapts business port methods +by method reference; the ports retain their business vocabulary and do not +inherit generic JDK function names. Receiver successes now use the same debug +logging path as submission, handling, and sending. diff --git a/README.md b/README.md index a7fd10a..c308db8 100644 --- a/README.md +++ b/README.md @@ -44,8 +44,9 @@ The modules have intentionally narrow responsibilities: - `adapters.messaging.postgresql` implements the transactional outbox. - `adapters.messaging.artemis` implements durable queues, topics, subscriptions, and polling with the Artemis Core client. -- `observability` decorates API and messaging boundaries with JDK logging, - concurrent counters, latency histograms, immutable snapshots, and JMX. +- `observability` adapts business ports through JDK functional interfaces and + decorates them with JDK logging, concurrent counters, latency histograms, + immutable snapshots, and JMX. - `app` is the only composition root. Nothing depends on it. See [ARCHITECTURE.md](ARCHITECTURE.md) for exact production module @@ -187,8 +188,8 @@ The project’s long-term target is a compact, executable example showing that: APIs; - push handlers and pull receivers can share minimal messaging ports; - durable retries can be explicit without embedding workflow logic in tables; -- logging and metrics can be composed around boundaries without annotations or - framework-managed proxies; and +- logging and metrics can be composed once around JDK functions, then adapted + to business ports without annotations or framework-managed proxies; and - new HTTP, database, broker, LocalStack, or test adapters can be added as new modules rather than edits to core rules. diff --git a/observability/src/main/java/dev/minimalaccounting/observability/Observability.java b/observability/src/main/java/dev/minimalaccounting/observability/Observability.java index 239cb23..c1c048d 100644 --- a/observability/src/main/java/dev/minimalaccounting/observability/Observability.java +++ b/observability/src/main/java/dev/minimalaccounting/observability/Observability.java @@ -1,34 +1,76 @@ package dev.minimalaccounting.observability; +import dev.minimalaccounting.ports.api.Transaction; import dev.minimalaccounting.ports.api.TransactionProcessor; import dev.minimalaccounting.ports.messaging.Handler; +import dev.minimalaccounting.ports.messaging.Received; import dev.minimalaccounting.ports.messaging.Receiver; import dev.minimalaccounting.ports.messaging.Sender; import lombok.NonNull; import lombok.RequiredArgsConstructor; +import java.util.List; +import java.util.function.Consumer; +import java.util.function.IntFunction; + @RequiredArgsConstructor public final class Observability { @NonNull private final MetricsRegistry metrics; - public TransactionProcessor processor(String name, TransactionProcessor processor) { - return new ObservedTransactionProcessor(name, processor, metrics, logger(name)); + public TransactionProcessor processor( + @NonNull String name, + @NonNull TransactionProcessor processor) { + Consumer observed = observer(name).consumer(name, processor::submit); + return observed::accept; + } + + public Handler handler(@NonNull String name, @NonNull Handler handler) { + Consumer observed = observer(name).consumer(name, handler::handle); + return observed::accept; } - public Handler handler(String name, Handler handler) { - return new ObservedHandler<>(name, handler, metrics, logger(name)); + public Sender sender(@NonNull String name, @NonNull Sender sender) { + Consumer observed = observer(name).consumer(name, sender::send); + return observed::accept; } - public Sender sender(String name, Sender sender) { - return new ObservedSender<>(name, sender, metrics, logger(name)); + public Receiver receiver(@NonNull String name, @NonNull Receiver receiver) { + OperationObserver observer = observer(name); + return new FunctionalReceiver<>( + observer.intFunction( + name + ".receive", + receiver::receive, + List::size), + observer.consumer(name + ".ack", receiver::ack), + observer.consumer(name + ".nack", receiver::nack)); } - public Receiver receiver(String name, Receiver receiver) { - return new ObservedReceiver<>(name, receiver, metrics, logger(name)); + private OperationObserver observer(String name) { + return new OperationObserver(metrics, logger(name)); } private static System.Logger logger(String name) { return System.getLogger("dev.minimalaccounting." + name); } + + private record FunctionalReceiver( + @NonNull IntFunction>> poll, + @NonNull Consumer acknowledge, + @NonNull Consumer reject) implements Receiver { + @Override + public List> receive(int limit) { + return poll.apply(limit); + } + + @Override + public void ack(String messageId) { + acknowledge.accept(messageId); + } + + @Override + public void nack(String messageId) { + reject.accept(messageId); + } + } } diff --git a/observability/src/main/java/dev/minimalaccounting/observability/ObservedHandler.java b/observability/src/main/java/dev/minimalaccounting/observability/ObservedHandler.java deleted file mode 100644 index 6f3482d..0000000 --- a/observability/src/main/java/dev/minimalaccounting/observability/ObservedHandler.java +++ /dev/null @@ -1,31 +0,0 @@ -package dev.minimalaccounting.observability; - -import dev.minimalaccounting.ports.messaging.Handler; -import lombok.NonNull; -import lombok.RequiredArgsConstructor; - -@RequiredArgsConstructor -public final class ObservedHandler implements Handler { - @NonNull - private final String name; - @NonNull - private final Handler delegate; - @NonNull - private final MetricsRegistry metrics; - @NonNull - private final System.Logger logger; - - @Override - public void handle(T message) { - long started = System.nanoTime(); - try { - delegate.handle(message); - metrics.success(name, System.nanoTime() - started, 1); - logger.log(System.Logger.Level.DEBUG, "{0} succeeded", name); - } catch (RuntimeException failure) { - metrics.failure(name, System.nanoTime() - started); - logger.log(System.Logger.Level.WARNING, "{0} failed", name); - throw failure; - } - } -} diff --git a/observability/src/main/java/dev/minimalaccounting/observability/ObservedReceiver.java b/observability/src/main/java/dev/minimalaccounting/observability/ObservedReceiver.java deleted file mode 100644 index 38bd236..0000000 --- a/observability/src/main/java/dev/minimalaccounting/observability/ObservedReceiver.java +++ /dev/null @@ -1,63 +0,0 @@ -package dev.minimalaccounting.observability; - -import dev.minimalaccounting.ports.messaging.Received; -import dev.minimalaccounting.ports.messaging.Receiver; -import lombok.NonNull; -import lombok.RequiredArgsConstructor; - -import java.util.List; - -@RequiredArgsConstructor -public final class ObservedReceiver implements Receiver { - @NonNull - private final String name; - @NonNull - private final Receiver delegate; - @NonNull - private final MetricsRegistry metrics; - @NonNull - private final System.Logger logger; - - @Override - public List> receive(int limit) { - long started = System.nanoTime(); - try { - List> messages = delegate.receive(limit); - metrics.success(name + ".receive", System.nanoTime() - started, messages.size()); - return messages; - } catch (RuntimeException failure) { - failed(name + ".receive", started, failure); - throw failure; - } - } - - @Override - public void ack(String messageId) { - disposition(name + ".ack", messageId, true); - } - - @Override - public void nack(String messageId) { - disposition(name + ".nack", messageId, false); - } - - private void disposition(String operation, String messageId, boolean acknowledge) { - long started = System.nanoTime(); - try { - if (acknowledge) { - delegate.ack(messageId); - } else { - delegate.nack(messageId); - } - metrics.success(operation, System.nanoTime() - started, 1); - } catch (RuntimeException failure) { - failed(operation, started, failure); - throw failure; - } - } - - private void failed(String operation, long started, RuntimeException failure) { - metrics.failure(operation, System.nanoTime() - started); - logger.log(System.Logger.Level.WARNING, "{0} failed", operation); - } -} diff --git a/observability/src/main/java/dev/minimalaccounting/observability/ObservedSender.java b/observability/src/main/java/dev/minimalaccounting/observability/ObservedSender.java deleted file mode 100644 index b719235..0000000 --- a/observability/src/main/java/dev/minimalaccounting/observability/ObservedSender.java +++ /dev/null @@ -1,31 +0,0 @@ -package dev.minimalaccounting.observability; - -import dev.minimalaccounting.ports.messaging.Sender; -import lombok.NonNull; -import lombok.RequiredArgsConstructor; - -@RequiredArgsConstructor -public final class ObservedSender implements Sender { - @NonNull - private final String name; - @NonNull - private final Sender delegate; - @NonNull - private final MetricsRegistry metrics; - @NonNull - private final System.Logger logger; - - @Override - public void send(T message) { - long started = System.nanoTime(); - try { - delegate.send(message); - metrics.success(name, System.nanoTime() - started, 1); - logger.log(System.Logger.Level.DEBUG, "{0} succeeded", name); - } catch (RuntimeException failure) { - metrics.failure(name, System.nanoTime() - started); - logger.log(System.Logger.Level.WARNING, "{0} failed", name); - throw failure; - } - } -} diff --git a/observability/src/main/java/dev/minimalaccounting/observability/ObservedTransactionProcessor.java b/observability/src/main/java/dev/minimalaccounting/observability/ObservedTransactionProcessor.java deleted file mode 100644 index a2895b6..0000000 --- a/observability/src/main/java/dev/minimalaccounting/observability/ObservedTransactionProcessor.java +++ /dev/null @@ -1,32 +0,0 @@ -package dev.minimalaccounting.observability; - -import dev.minimalaccounting.ports.api.Transaction; -import dev.minimalaccounting.ports.api.TransactionProcessor; -import lombok.NonNull; -import lombok.RequiredArgsConstructor; - -@RequiredArgsConstructor -public final class ObservedTransactionProcessor implements TransactionProcessor { - @NonNull - private final String name; - @NonNull - private final TransactionProcessor delegate; - @NonNull - private final MetricsRegistry metrics; - @NonNull - private final System.Logger logger; - - @Override - public void submit(Transaction transaction) { - long started = System.nanoTime(); - try { - delegate.submit(transaction); - metrics.success(name, System.nanoTime() - started, 1); - logger.log(System.Logger.Level.DEBUG, "{0} succeeded", name); - } catch (RuntimeException failure) { - metrics.failure(name, System.nanoTime() - started); - logger.log(System.Logger.Level.WARNING, "{0} failed", name); - throw failure; - } - } -} diff --git a/observability/src/main/java/dev/minimalaccounting/observability/OperationObserver.java b/observability/src/main/java/dev/minimalaccounting/observability/OperationObserver.java new file mode 100644 index 0000000..9853043 --- /dev/null +++ b/observability/src/main/java/dev/minimalaccounting/observability/OperationObserver.java @@ -0,0 +1,56 @@ +package dev.minimalaccounting.observability; + +import lombok.NonNull; +import lombok.RequiredArgsConstructor; + +import java.util.function.Consumer; +import java.util.function.IntFunction; +import java.util.function.Supplier; +import java.util.function.ToLongFunction; + +@RequiredArgsConstructor +final class OperationObserver { + @NonNull + private final MetricsRegistry metrics; + @NonNull + private final System.Logger logger; + + Consumer consumer( + @NonNull String operation, + @NonNull Consumer delegate) { + return argument -> observe(operation, () -> { + delegate.accept(argument); + return null; + }, ignored -> 1); + } + + IntFunction intFunction( + @NonNull String operation, + @NonNull IntFunction delegate, + @NonNull ToLongFunction itemCount) { + return argument -> observe( + operation, + () -> delegate.apply(argument), + itemCount); + } + + private T observe( + String operation, + Supplier invocation, + ToLongFunction itemCount) { + long started = System.nanoTime(); + try { + T result = invocation.get(); + metrics.success( + operation, + System.nanoTime() - started, + itemCount.applyAsLong(result)); + logger.log(System.Logger.Level.DEBUG, "{0} succeeded", operation); + return result; + } catch (RuntimeException failure) { + metrics.failure(operation, System.nanoTime() - started); + logger.log(System.Logger.Level.WARNING, "{0} failed", operation); + throw failure; + } + } +} diff --git a/observability/src/test/java/dev/minimalaccounting/observability/ObservedMessagingTest.java b/observability/src/test/java/dev/minimalaccounting/observability/ObservabilityTest.java similarity index 56% rename from observability/src/test/java/dev/minimalaccounting/observability/ObservedMessagingTest.java rename to observability/src/test/java/dev/minimalaccounting/observability/ObservabilityTest.java index e8aba47..f71ed0d 100644 --- a/observability/src/test/java/dev/minimalaccounting/observability/ObservedMessagingTest.java +++ b/observability/src/test/java/dev/minimalaccounting/observability/ObservabilityTest.java @@ -6,44 +6,57 @@ import java.time.Clock; import java.util.List; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertThrows; -class ObservedMessagingTest { +class ObservabilityTest { @Test - void senderRecordsSuccessAndPreservesTheMessage() { + void adaptsSingleArgumentPortsThroughObservedConsumers() { var metrics = new MetricsRegistry(Clock.systemUTC()); - var sent = new java.util.concurrent.atomic.AtomicReference(); - var sender = new ObservedSender( - "sender", sent::set, metrics, System.getLogger("test")); - - sender.send("message"); - - assertEquals("message", sent.get()); - assertEquals(1, metrics.snapshot().operations().get("sender").successes()); + var observability = new Observability(metrics); + var submissions = new AtomicInteger(); + var handled = new AtomicReference(); + var sent = new AtomicReference(); + + observability.processor("transaction.submit", ignored -> submissions.incrementAndGet()) + .submit(null); + observability.handler("command.handle", handled::set).handle("command"); + observability.sender("event.send", sent::set).send("event"); + + assertEquals(1, submissions.get()); + assertEquals("command", handled.get()); + assertEquals("event", sent.get()); + var operations = metrics.snapshot().operations(); + assertEquals(1, operations.get("transaction.submit").successes()); + assertEquals(1, operations.get("command.handle").successes()); + assertEquals(1, operations.get("event.send").successes()); } @Test - void receiverCountsMessagesAndDispositions() { + void receiverAdaptsPollingAndBothDispositions() { var metrics = new MetricsRegistry(Clock.systemUTC()); var delegate = new RecordingReceiver(); - var receiver = new ObservedReceiver<>( - "queue", delegate, metrics, System.getLogger("test")); + var receiver = new Observability(metrics).receiver("queue", delegate); assertEquals(2, receiver.receive(2).size()); receiver.ack("1"); receiver.nack("2"); - var operations = metrics.snapshot().operations(); - assertEquals(2, operations.get("queue.receive").items()); + assertEquals(2, delegate.limit); assertEquals("1", delegate.acked); assertEquals("2", delegate.nacked); + var operations = metrics.snapshot().operations(); + assertEquals(2, operations.get("queue.receive").items()); + assertEquals(1, operations.get("queue.ack").successes()); + assertEquals(1, operations.get("queue.nack").successes()); } @Test - void receiverRethrowsDispositionFailure() { + void receiverRethrowsTheSameDispositionFailure() { var metrics = new MetricsRegistry(Clock.systemUTC()); var expected = new IllegalArgumentException("unknown"); Receiver delegate = new RecordingReceiver() { @@ -52,8 +65,7 @@ public void ack(String messageId) { throw expected; } }; - var receiver = new ObservedReceiver<>( - "queue", delegate, metrics, System.getLogger("test")); + var receiver = new Observability(metrics).receiver("queue", delegate); var thrown = assertThrows(IllegalArgumentException.class, () -> receiver.ack("missing")); @@ -63,11 +75,13 @@ public void ack(String messageId) { } private static class RecordingReceiver implements Receiver { + private int limit; private String acked; private String nacked; @Override public List> receive(int limit) { + this.limit = limit; return List.of(new Received<>("1", "first"), new Received<>("2", "second")); } diff --git a/observability/src/test/java/dev/minimalaccounting/observability/ObservedHandlerTest.java b/observability/src/test/java/dev/minimalaccounting/observability/ObservedHandlerTest.java deleted file mode 100644 index c57f7d7..0000000 --- a/observability/src/test/java/dev/minimalaccounting/observability/ObservedHandlerTest.java +++ /dev/null @@ -1,43 +0,0 @@ -package dev.minimalaccounting.observability; - -import org.junit.jupiter.api.Test; - -import java.time.Clock; - -import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertSame; -import static org.junit.jupiter.api.Assertions.assertThrows; - -class ObservedHandlerTest { - @Test - void recordsSuccessWithoutChangingTheMessage() { - var metrics = new MetricsRegistry(Clock.systemUTC()); - var received = new java.util.concurrent.atomic.AtomicReference(); - var handler = new ObservedHandler( - "handler", received::set, metrics, System.getLogger("test")); - - handler.handle("message"); - - assertEquals("message", received.get()); - assertEquals(1, metrics.snapshot().operations().get("handler").successes()); - } - - @Test - void recordsAndRethrowsTheSameRuntimeFailure() { - var metrics = new MetricsRegistry(Clock.systemUTC()); - var expected = new IllegalStateException("failure"); - var handler = new ObservedHandler( - "handler", - ignored -> { - throw expected; - }, - metrics, - System.getLogger("test")); - - var thrown = assertThrows(IllegalStateException.class, - () -> handler.handle("message")); - - assertSame(expected, thrown); - assertEquals(1, metrics.snapshot().operations().get("handler").failures()); - } -} diff --git a/observability/src/test/java/dev/minimalaccounting/observability/ObservedTransactionProcessorTest.java b/observability/src/test/java/dev/minimalaccounting/observability/ObservedTransactionProcessorTest.java deleted file mode 100644 index 54f41d4..0000000 --- a/observability/src/test/java/dev/minimalaccounting/observability/ObservedTransactionProcessorTest.java +++ /dev/null @@ -1,29 +0,0 @@ -package dev.minimalaccounting.observability; - -import dev.minimalaccounting.ports.api.TransactionProcessor; -import org.junit.jupiter.api.Test; - -import java.time.Clock; - -import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertSame; -import static org.junit.jupiter.api.Assertions.assertThrows; - -class ObservedTransactionProcessorTest { - @Test - void recordsAndRethrowsSubmissionFailure() { - var metrics = new MetricsRegistry(Clock.systemUTC()); - var expected = new IllegalArgumentException("invalid"); - TransactionProcessor delegate = ignored -> { - throw expected; - }; - var processor = new ObservedTransactionProcessor( - "submit", delegate, metrics, System.getLogger("test")); - - var thrown = assertThrows(IllegalArgumentException.class, - () -> processor.submit(null)); - - assertSame(expected, thrown); - assertEquals(1, metrics.snapshot().operations().get("submit").failures()); - } -} diff --git a/observability/src/test/java/dev/minimalaccounting/observability/OperationObserverTest.java b/observability/src/test/java/dev/minimalaccounting/observability/OperationObserverTest.java new file mode 100644 index 0000000..aacf51e --- /dev/null +++ b/observability/src/test/java/dev/minimalaccounting/observability/OperationObserverTest.java @@ -0,0 +1,100 @@ +package dev.minimalaccounting.observability; + +import org.junit.jupiter.api.Test; + +import java.time.Clock; +import java.util.ArrayList; +import java.util.List; +import java.util.ResourceBundle; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertThrows; + +class OperationObserverTest { + @Test + void consumerRecordsSuccessAndPreservesTheArgument() { + var metrics = new MetricsRegistry(Clock.systemUTC()); + var logger = new RecordingLogger(); + var received = new AtomicReference(); + var consumer = new OperationObserver(metrics, logger) + .consumer("handler", received::set); + + consumer.accept("message"); + + assertEquals("message", received.get()); + assertEquals(1, metrics.snapshot().operations().get("handler").successes()); + assertEquals(List.of(System.Logger.Level.DEBUG), logger.levels); + } + + @Test + void consumerRecordsFailureAndRethrowsTheSameException() { + var metrics = new MetricsRegistry(Clock.systemUTC()); + var logger = new RecordingLogger(); + var expected = new IllegalStateException("failure"); + var consumer = new OperationObserver(metrics, logger) + .consumer("handler", ignored -> { + throw expected; + }); + + var thrown = assertThrows(IllegalStateException.class, + () -> consumer.accept("message")); + + assertSame(expected, thrown); + assertEquals(1, metrics.snapshot().operations().get("handler").failures()); + assertEquals(List.of(System.Logger.Level.WARNING), logger.levels); + } + + @Test + void intFunctionPreservesTheResultAndCountsItsItems() { + var metrics = new MetricsRegistry(Clock.systemUTC()); + var logger = new RecordingLogger(); + var expected = List.of("first", "second"); + var function = new OperationObserver(metrics, logger) + .intFunction("queue.receive", limit -> { + assertEquals(2, limit); + return expected; + }, List::size); + + var result = function.apply(2); + + assertSame(expected, result); + var snapshot = metrics.snapshot().operations().get("queue.receive"); + assertEquals(1, snapshot.successes()); + assertEquals(2, snapshot.items()); + assertEquals(List.of(System.Logger.Level.DEBUG), logger.levels); + } + + private static final class RecordingLogger implements System.Logger { + private final List levels = new ArrayList<>(); + + @Override + public String getName() { + return "test"; + } + + @Override + public boolean isLoggable(Level level) { + return true; + } + + @Override + public void log( + Level level, + ResourceBundle bundle, + String message, + Throwable thrown) { + levels.add(level); + } + + @Override + public void log( + Level level, + ResourceBundle bundle, + String format, + Object... parameters) { + levels.add(level); + } + } +}