diff --git a/purge-payments.hurl b/purge-payments.hurl new file mode 100644 index 0000000..8ed518a --- /dev/null +++ b/purge-payments.hurl @@ -0,0 +1,2 @@ +POST http://localhost:8002/admin/purge-payments +X-Rinha-Token: 123 diff --git a/request-get-all.hurl b/request-get-all.hurl new file mode 100644 index 0000000..7c143ba --- /dev/null +++ b/request-get-all.hurl @@ -0,0 +1 @@ +GET http://localhost:8000/payments-summary diff --git a/set-fail.hurl b/set-fail.hurl index 446db7d..b64c948 100644 --- a/set-fail.hurl +++ b/set-fail.hurl @@ -1,5 +1,5 @@ -PUT http://localhost:8001/admin/configurations/failure +PUT http://localhost:8002/admin/configurations/failure X-Rinha-Token: 123 { - "failure": true + "failure": false } diff --git a/src/integrations/fallback.gleam b/src/integrations/fallback.gleam deleted file mode 100644 index eb66388..0000000 --- a/src/integrations/fallback.gleam +++ /dev/null @@ -1,19 +0,0 @@ -import gleam/http -import gleam/http/request -import gleam/httpc -import gleam/result -import integrations/default - -pub fn fallback_provider_send_request(body: String) { - let assert Ok(request) = request.to("http://localhost:8002/payments") - - use response <- result.try( - request - |> request.set_header("content-type", "application/json") - |> request.set_body(body) - |> request.set_method(http.Post) - |> httpc.send, - ) - - default.parse_http_response(response) -} diff --git a/src/integrations/default.gleam b/src/integrations/provider.gleam similarity index 82% rename from src/integrations/default.gleam rename to src/integrations/provider.gleam index 10fc71a..75238fe 100644 --- a/src/integrations/default.gleam +++ b/src/integrations/provider.gleam @@ -7,6 +7,10 @@ import gleam/json import gleam/result import model +pub type ProviderConfig { + ProviderConfig(url: String) +} + pub fn create_body(body: model.PaymentRequest) -> String { json.object([ #("correlationId", json.string(body.correlation_id)), @@ -16,8 +20,8 @@ pub fn create_body(body: model.PaymentRequest) -> String { |> json.to_string } -pub fn default_provider_send_request(body: String) { - let assert Ok(request) = request.to("http://localhost:8001/payments") +pub fn send_request(provider: ProviderConfig, body: String) { + let assert Ok(request) = request.to(provider.url <> "/payments") use response <- result.try( request diff --git a/src/model.gleam b/src/model.gleam index 26e1306..8fd37b2 100644 --- a/src/model.gleam +++ b/src/model.gleam @@ -1,5 +1,41 @@ import birl +import gleam/dict +import gleam/dynamic/decode +import gleam/json +import gleam/result pub type PaymentRequest { PaymentRequest(correlation_id: String, amount: Float, requested_at: birl.Time) } + +pub fn to_dict(payment: PaymentRequest) -> dict.Dict(String, String) { + dict.new() + |> dict.insert(payment.correlation_id, payment |> to_json |> json.to_string) +} + +pub fn to_json(payment: PaymentRequest) -> json.Json { + [ + #("correlationId", json.string(payment.correlation_id)), + #("amount", json.float(payment.amount)), + #("requestedAt", json.string(payment.requested_at |> birl.to_iso8601)), + ] + |> json.object +} + +pub fn from_json_string(payment: String) { + let parse = { + use amount <- decode.field("amount", decode.float) + use correlation_id <- decode.field("correlationId", decode.string) + use requested_at <- decode.field("requestedAt", decode.string) + + decode.success(PaymentRequest( + amount: amount, + correlation_id: correlation_id, + requested_at: requested_at + |> birl.parse + |> result.unwrap(birl.now()), + )) + } + + json.parse(payment, parse) +} diff --git a/src/processor.gleam b/src/processor.gleam index b4505c9..9059324 100644 --- a/src/processor.gleam +++ b/src/processor.gleam @@ -1,14 +1,9 @@ -import birl -import gleam/dynamic/decode import gleam/erlang/process -import gleam/json import gleam/option import gleam/otp/actor import gleam/otp/supervision -import gleam/result -import integrations/default -import integrations/fallback -import model.{PaymentRequest} +import integrations/provider +import model import redis import valkyrie @@ -16,6 +11,8 @@ pub type Processor { Processor( redis_conn: valkyrie.Connection, name: option.Option(process.Name(Message)), + connections: Int, + providers: List(provider.ProviderConfig), ) } @@ -23,14 +20,25 @@ pub type Message { ServerTick } -pub fn new(conn: valkyrie.Connection) { - Processor(redis_conn: conn, name: option.None) +pub fn new(conn: valkyrie.Connection) -> Processor { + Processor(redis_conn: conn, name: option.None, connections: 1, providers: []) } -pub fn named(processor: Processor, name: process.Name(Message)) { +pub fn named(processor: Processor, name: process.Name(Message)) -> Processor { Processor(..processor, name: option.Some(name)) } +pub fn connections(processor: Processor, connections: Int) -> Processor { + Processor(..processor, connections: connections) +} + +pub fn providers( + processor: Processor, + providers: List(provider.ProviderConfig), +) -> Processor { + Processor(..processor, providers: providers) +} + pub fn start(processor: Processor) { let ac = processor @@ -49,30 +57,12 @@ pub fn supervised(processor: Processor) { } fn handle_message(state: Processor, message: Message) { - echo "inside on message" case message { ServerTick -> { case redis.read_queue_payments(state.redis_conn) { "" -> actor.continue(state) data -> { - let assert Ok(message) = parse_data_redis(data) - - let body_to_send = message |> default.create_body - - case default.default_provider_send_request(body_to_send) { - Ok(_) -> actor.continue(state) - Error(_) -> { - echo "failed to make request" - case fallback.fallback_provider_send_request(body_to_send) { - Ok(_) -> actor.continue(state) - Error(_) -> { - echo "failed to make request" - actor.continue(state) - } - } - } - } - + process.spawn(fn() { integrate_data(state, data) }) actor.continue(state) } } @@ -81,27 +71,45 @@ fn handle_message(state: Processor, message: Message) { } pub fn loop_worker(subject: process.Subject(Message)) { - echo "inside pool" process.send(subject, ServerTick) process.sleep(2000) loop_worker(subject) } -fn parse_data_redis(message: String) { - echo message - let parse = { - use amount <- decode.field("amount", decode.float) - use correlation_id <- decode.field("correlationId", decode.string) - use requested_at <- decode.field("requestedAt", decode.string) - - decode.success(PaymentRequest( - amount: amount, - correlation_id: correlation_id, - requested_at: requested_at - |> birl.parse - |> result.unwrap(birl.now()), - )) +fn integrations( + providers: List(provider.ProviderConfig), + body: String, +) -> Result(Bool, Bool) { + case providers { + [] -> Error(False) + [provider, ..rest] -> { + echo "Trying provider -> " <> provider.url + case provider.send_request(provider, body) { + Ok(_) -> Ok(True) + _ -> integrations(rest, body) + } + } } +} - json.parse(message, parse) +fn integrate_data(processor: Processor, data: String) { + echo "Trying to integrate " <> data + let assert Ok(message) = model.from_json_string(data) + + let body_to_send = message |> provider.create_body + + case integrations(processor.providers, body_to_send) { + Error(_) -> { + echo "Failed, reenqueuing it" + + processor.redis_conn + |> redis.enqueue_payments([data]) + } + Ok(_) -> { + echo "Success, saving it" + + processor.redis_conn + |> redis.save_data(message |> model.to_dict) + } + } } diff --git a/src/redis.gleam b/src/redis.gleam index e15507a..1a4f251 100644 --- a/src/redis.gleam +++ b/src/redis.gleam @@ -1,15 +1,19 @@ +import gleam/dict import gleam/erlang/process import gleam/option import valkyrie const default_timeout = 1000 -pub fn create_supervised_pool() { +const default_key = "payments" + +pub fn create_supervised_pool(host: String) { let name = process.new_name("connection_pool") #( name, valkyrie.default_config() + |> valkyrie.host(host) |> valkyrie.supervised_pool( size: 10_000, name: option.Some(name), @@ -18,8 +22,9 @@ pub fn create_supervised_pool() { ) } -pub fn enqueue_payments(body: List(String), conn: valkyrie.Connection) { - valkyrie.lpush(conn, "payments_created", body, default_timeout) +pub fn enqueue_payments(conn: valkyrie.Connection, body: List(String)) { + conn + |> valkyrie.lpush("payments_created", body, default_timeout) } pub fn read_queue_payments(conn: valkyrie.Connection) -> String { @@ -33,3 +38,13 @@ pub fn read_queue_payments(conn: valkyrie.Connection) -> String { Error(_) -> "" } } + +pub fn save_data(conn: valkyrie.Connection, body: dict.Dict(String, String)) { + conn + |> valkyrie.hset(default_key, body, default_timeout) +} + +pub fn get_all_saved_data(conn: valkyrie.Connection) { + conn + |> valkyrie.hgetall(default_key, default_timeout) +} diff --git a/src/rinha_2025.gleam b/src/rinha_2025.gleam index f2d2a2e..fcf7724 100644 --- a/src/rinha_2025.gleam +++ b/src/rinha_2025.gleam @@ -1,20 +1,36 @@ +import gleam/bool import gleam/erlang/process +import gleam/list import gleam/otp/static_supervisor as supervisor import gleam/otp/supervision +import gleam/string +import integrations/provider import processor import redis +import util import valkyrie import web/server import web/web pub fn main() -> Nil { - let #(valkey_pool_name, valkey_pool) = redis.create_supervised_pool() + let redis_host = util.get_env_var("REDIS_CONN", "localhost") + let providers_env = util.get_env_var("PROVIDERS", "") + let providers = + providers_env + |> string.split(",") + |> list.filter(fn(url) { "" != url }) + |> list.map(fn(url) { provider.ProviderConfig(url: url) }) + let has_providers = list.length(providers) > 0 + + let #(valkey_pool_name, valkey_pool) = + redis.create_supervised_pool(redis_host) let valky = valkyrie.named_connection(valkey_pool_name) let worker_name = process.new_name("worker_pool") let worker_pool_supervised = processor.new(valky) |> processor.named(worker_name) + |> processor.providers(providers) |> processor.supervised let ctx = server.Context(valkye_conn: valky) @@ -23,10 +39,12 @@ pub fn main() -> Nil { supervisor.new(supervisor.OneForOne) |> supervisor.add(valkey_pool) |> supervisor.add(web.create_server_supervised(ctx)) - |> supervisor.add(worker_pool_supervised) + |> create_supervisor_with_processor(worker_pool_supervised, has_providers) |> supervisor.start process.spawn(fn() { + use <- bool.guard(when: !has_providers, return: Ok(Nil)) + worker_name |> process.named_subject |> processor.loop_worker @@ -34,3 +52,14 @@ pub fn main() -> Nil { process.sleep_forever() } + +fn create_supervisor_with_processor( + manager: supervisor.Builder, + processor: supervision.ChildSpecification(process.Subject(processor.Message)), + has_processor: Bool, +) { + use <- bool.guard(when: !has_processor, return: manager) + + manager + |> supervisor.add(processor) +} diff --git a/src/util.gleam b/src/util.gleam new file mode 100644 index 0000000..97ee048 --- /dev/null +++ b/src/util.gleam @@ -0,0 +1,8 @@ +import envoy + +pub fn get_env_var(env: String, default: String) -> String { + case envoy.get(env) { + Error(_) -> default + Ok(val) -> val + } +} diff --git a/src/web/controllers/payment_controller.gleam b/src/web/controllers/payment_controller.gleam index 1b40978..736964a 100644 --- a/src/web/controllers/payment_controller.gleam +++ b/src/web/controllers/payment_controller.gleam @@ -1,13 +1,14 @@ import birl +import gleam/dict import gleam/dynamic import gleam/dynamic/decode -import gleam/erlang/process +import gleam/http/response import gleam/json import gleam/list -import gleam/otp/actor +import gleam/string_tree import model.{type PaymentRequest, PaymentRequest} -import processor import redis +import valkyrie/resp import web/server import wisp @@ -31,21 +32,48 @@ fn decode_payment_body( } } -pub fn handle_payment_post(req: wisp.Request, ctx: server.Context) { +pub fn handle_payment_post( + req: wisp.Request, + ctx: server.Context, +) -> response.Response(wisp.Body) { use json <- wisp.require_json(req) use body <- decode_payment_body(json) - let assert Ok(_) = - echo json.object([ - #("correlationId", json.string(body.correlation_id)), - #("amount", json.float(body.amount)), - #("requestedAt", json.string(birl.now() |> birl.to_iso8601)), - ]) - |> json.to_string - |> list.wrap - |> redis.enqueue_payments(ctx.valkye_conn) + let data_to_insert = + body + |> model.to_json + |> json.to_string + |> list.wrap - // actor.send(ctx.worker_subject, processor.Process(body)) + let assert Ok(_) = + echo ctx.valkye_conn + |> redis.enqueue_payments(data_to_insert) wisp.no_content() } + +pub fn get_all_payments( + _req: wisp.Request, + ctx: server.Context, +) -> response.Response(wisp.Body) { + case redis.get_all_saved_data(ctx.valkye_conn) { + Ok(data) -> { + let response = + data + |> dict.values + |> list.map(fn(value) { + let assert resp.BulkString(v) = value + let assert Ok(json_data) = model.from_json_string(v) + json_data + }) + |> json.array(of: model.to_json) + |> json.to_string_tree + + wisp.ok() + |> wisp.json_body(response) + } + Error(_) -> + wisp.no_content() + |> wisp.json_body(string_tree.from_string("[]")) + } +} diff --git a/src/web/router.gleam b/src/web/router.gleam index 2fcbc36..5aa6640 100644 --- a/src/web/router.gleam +++ b/src/web/router.gleam @@ -14,7 +14,7 @@ pub fn handle_request(req: wisp.Request, ctx: server.Context) -> wisp.Response { } ["payments-summary"] -> { use <- wisp.require_method(req, http.Get) - todo + payment_controller.get_all_payments(req, ctx) } _ -> wisp.not_found() } diff --git a/src/web/server.gleam b/src/web/server.gleam index 6b79280..e46a8ca 100644 --- a/src/web/server.gleam +++ b/src/web/server.gleam @@ -1,10 +1,5 @@ -import gleam/erlang/process -import processor import valkyrie pub type Context { - Context( - valkye_conn: valkyrie.Connection, - // worker_subject: process.Subject(processor.Message), - ) + Context(valkye_conn: valkyrie.Connection) }