This repository has been archived by the owner on May 11, 2023. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
20 changed files
with
314 additions
and
84 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,3 +1,4 @@ | ||
--format documentation | ||
--format progress | ||
--order random | ||
--color | ||
--require spec_helper |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -41,3 +41,6 @@ Style/WordArray: | |
|
||
Style/ClassAndModuleChildren: | ||
EnforcedStyle: compact | ||
|
||
RSpec/MultipleMemoizedHelpers: | ||
Max: 10 |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -3,6 +3,7 @@ | |
source "https://rubygems.org" | ||
|
||
gem "async" | ||
gem "async-http" | ||
gem "nats-pure" | ||
|
||
gem "faraday" | ||
|
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,19 @@ | ||
# frozen_string_literal: true | ||
|
||
class NatsStreamer::ActiveTasks | ||
include NatsStreamer::Helpers | ||
|
||
def async | ||
tasks << Async do |task| | ||
task.yield | ||
yield(task) | ||
tasks.delete(task) | ||
end | ||
end | ||
|
||
def wait = tasks.each(&:wait) | ||
|
||
private | ||
|
||
memoize def tasks = [] | ||
end |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,26 +1,50 @@ | ||
# frozen_string_literal: true | ||
|
||
class NatsStreamer::Application | ||
extend Dry::Initializer | ||
include NatsStreamer::Helpers | ||
|
||
include NatsStreamer::Logger | ||
include Memery | ||
|
||
option :config | ||
option :config, type: T.Instance(NatsStreamer::Config) | ||
|
||
def run | ||
config.streams.each do |name, subjects| | ||
NatsStreamer::Stream.new(jsm:, name:, subjects:).run | ||
end | ||
metrics_store.run | ||
metrics_server.run | ||
wait | ||
ensure | ||
stop | ||
wait | ||
|
||
client.close | ||
end | ||
|
||
def stop | ||
metrics_server.stop | ||
streams.each(&:stop) | ||
metrics_store.stop | ||
end | ||
|
||
Async::Notification.new.wait # wait forever | ||
def wait | ||
metrics_server.wait | ||
streams.each(&:wait) | ||
metrics_store.wait | ||
end | ||
|
||
private | ||
|
||
memoize def jsm | ||
memoize def jsm = client.jsm | ||
|
||
# TODO: move port to config | ||
memoize def metrics_server = NatsStreamer::Metrics::Server.new(port: 9294, metrics_store:) | ||
memoize def metrics_store = NatsStreamer::Metrics::Store.new | ||
|
||
memoize def streams | ||
config.streams.map do |name, subjects| | ||
NatsStreamer::Stream.new(jsm:, name:, subjects:, metrics_store:).tap(&:run) | ||
end | ||
end | ||
|
||
memoize def client | ||
NATS.connect(config.server_url).tap do | ||
info { "Connected to #{config.server_url}" } | ||
end.jsm | ||
end | ||
end | ||
end |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,24 +1,17 @@ | ||
# frozen_string_literal: true | ||
|
||
class NatsStreamer::Config < Dry::Struct | ||
Types = Dry.Types | ||
T = Dry.Types | ||
|
||
class Subscriber < Dry::Struct | ||
attribute :name, Types::String # Must be uniq | ||
attribute :url, Types::String | ||
attribute :params, Types::Hash.map(Types::Coercible::String, Types::String).default({}.freeze) | ||
attribute :name, T::String # Must be uniq | ||
attribute :url, T::String | ||
attribute :params, T::Hash.map(T::Coercible::String, T::String).default({}.freeze) | ||
end | ||
|
||
attribute :server_url, Types::Coercible::String | ||
Events = T::Hash.map(T::Coercible::String, T::Array.of(Subscriber)) | ||
Subjects = T::Hash.map(T::Coercible::String, Events) | ||
|
||
attribute :streams, Types::Hash.map( | ||
Types::Coercible::String, | ||
Types::Hash.map( | ||
Types::Coercible::String, | ||
Types::Hash.map( | ||
Types::Coercible::String, | ||
Types::Array.of(Subscriber) | ||
) | ||
) | ||
) | ||
attribute :server_url, T::Coercible::String | ||
attribute :streams, T::Hash.map(T::Coercible::String, Subjects) | ||
end |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,16 +1,20 @@ | ||
# frozen_string_literal: true | ||
|
||
class NatsStreamer::Deliverer | ||
extend Dry::Initializer | ||
include NatsStreamer::Helpers | ||
|
||
include NatsStreamer::Logger | ||
include Memery | ||
option :subscriber, type: T.Instance(NatsStreamer::Config::Subscriber) | ||
option :metrics_store, type: T.Instance(NatsStreamer::Metrics::Store) | ||
|
||
option :subscriber | ||
|
||
def deliver(**) | ||
info_measure(-> { connection.post(".", **) }) { "Event delivered to #{subscriber.name}: #{_1.round(2)}s" } | ||
def deliver(**params) | ||
name = params.dig(:event, :name) | ||
info_measure(-> { "Event #{name.inspect} delivered to #{subscriber.name.inspect}: #{_1.round(2)}s" }) do | ||
connection.post(".", **params) | ||
end | ||
metrics_store.inc(:delivered, subscriber: subscriber.name) | ||
end | ||
|
||
private | ||
|
||
memoize def connection = NatsStreamer::Connection.build(subscriber.url) | ||
end |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,11 @@ | ||
# frozen_string_literal: true | ||
|
||
module NatsStreamer::Helpers | ||
T = Dry.Types | ||
|
||
def self.included(base) | ||
base.extend(Dry::Initializer) | ||
base.include(NatsStreamer::Logger) | ||
base.include(Memery) | ||
end | ||
end |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,45 +1,38 @@ | ||
# frozen_string_literal: true | ||
|
||
class NatsStreamer::Listener | ||
extend Dry::Initializer | ||
include NatsStreamer::Helpers | ||
|
||
include NatsStreamer::Logger | ||
include Memery | ||
|
||
# TODO: how to propely unsubscribe? | ||
option :jsm | ||
option :subject | ||
option :subscriber | ||
option :jsm, type: T.Instance(NATS::JetStream) | ||
option :subject, type: T::Coercible::String | ||
option :subscriber, type: T.Instance(NatsStreamer::Config::Subscriber) | ||
option :metrics_store, type: T.Instance(NatsStreamer::Metrics::Store) | ||
|
||
def run | ||
info { "Subscribing to #{subject}" } | ||
|
||
pull do |msg| | ||
Async { handle_message(msg) } | ||
end | ||
puller.pull { handle_message(_1) } | ||
wait | ||
end | ||
|
||
private | ||
|
||
memoize def durable = "nats-streamer-subscriber-#{subscriber.name}" | ||
memoize def params_renderer = NatsStreamer::ParamsRenderer.new(subject:, params: subscriber.params) | ||
memoize def deliverer(subscriber) = NatsStreamer::Deliverer.new(subscriber:) | ||
|
||
def pull(&) | ||
psub = @jsm.pull_subscribe(subject, "nats-streamer-subscriber-#{subscriber.name}") | ||
memoize def active_tasks = NatsStreamer::ActiveTasks.new | ||
memoize def puller = NatsStreamer::Puller.new(jsm:, subject:, durable:) | ||
memoize def deliverer = NatsStreamer::Deliverer.new(subscriber:, metrics_store:) | ||
|
||
loop do | ||
psub.fetch(1).each(&) | ||
rescue NATS::IO::Timeout | ||
debug { "Pulling timeout, retrying..." } | ||
end | ||
end | ||
def wait = active_tasks.wait | ||
|
||
def handle_message(msg) | ||
event = JSON.parse(msg.data, symbolize_names: true) | ||
debug { "subject=#{subject}, event=#{event}" } | ||
active_tasks.async do | ||
event = JSON.parse(msg.data, symbolize_names: true) | ||
debug { "subject=#{subject}, event=#{event}" } | ||
|
||
deliverer(subscriber).deliver(**params_renderer.render(event)) | ||
deliverer.deliver(**params_renderer.render(event)) | ||
|
||
msg.ack | ||
msg.ack | ||
end | ||
end | ||
end |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,39 @@ | ||
# frozen_string_literal: true | ||
|
||
class NatsStreamer::Metrics::Server | ||
include NatsStreamer::Helpers | ||
|
||
PATHS = ["/metrics", "/metrics/"].freeze | ||
|
||
NOT_FOUND = Protocol::HTTP::Response[404, {}, ["Not found"]].freeze | ||
|
||
option :port, type: T::Integer | ||
option :metrics_store, type: T.Instance(NatsStreamer::Metrics::Store) | ||
|
||
def run = @task = Async { Async::HTTP::Server.new(method(:app), endpoint).run } | ||
|
||
def stop = @task.stop | ||
def wait = @task.wait | ||
|
||
private | ||
|
||
memoize def endpoint = Async::HTTP::Endpoint.parse("http://127.0.0.1:#{port}") | ||
|
||
def app(request) | ||
return NOT_FOUND unless PATHS.include?(request.path) | ||
|
||
Protocol::HTTP::Response[200, {}, serialize_metrics] | ||
end | ||
|
||
def serialize_metrics | ||
metrics_store.metrics.map do |value| | ||
"#{metric_name(value)}{#{metric_tags(value)}} #{value[:value]}" | ||
end.join("\n") | ||
end | ||
|
||
def metric_name(value, unit = "total") | ||
"nats_streamer_#{value[:name]}_#{unit}" | ||
end | ||
|
||
def metric_tags(value) = value[:tags].map { |tag, tag_value| "#{tag}=#{tag_value.to_s.inspect}" }.join(",") | ||
end |
Oops, something went wrong.