diff --git a/ruby_event_store-browser/lib/ruby_event_store/browser.rb b/ruby_event_store-browser/lib/ruby_event_store/browser.rb index c608df9439..763ad47e8e 100644 --- a/ruby_event_store-browser/lib/ruby_event_store/browser.rb +++ b/ruby_event_store-browser/lib/ruby_event_store/browser.rb @@ -22,6 +22,7 @@ def self.fingerprint(name) end require_relative "browser/get_events_from_stream" +require_relative "browser/get_events_from_streams" require_relative "browser/urls" require_relative "browser/router" require_relative "browser/renderer" diff --git a/ruby_event_store-browser/lib/ruby_event_store/browser/app.rb b/ruby_event_store-browser/lib/ruby_event_store/browser/app.rb index 33930177df..31e5257d30 100644 --- a/ruby_event_store-browser/lib/ruby_event_store/browser/app.rb +++ b/ruby_event_store-browser/lib/ruby_event_store/browser/app.rb @@ -162,6 +162,41 @@ def call(env) ) end + router.add_route("GET", "/swimlane") do |params, urls| + stream_names, sort = swimlane_params(params) + reader = GetEventsFromStreams.new(event_store: event_store, stream_names: stream_names, sort: sort) + html render( + "swimlane/show", + urls: urls, + stream_names: stream_names, + events: reader.events, + sort: sort, + more_url: (urls.swimlane_more_url(stream_names, reader.next_cursor, sort) if reader.more?), + ) + end + + router.add_route("GET", "/swimlane/more") do |params, urls| + stream_names, sort = swimlane_params(params) + reader = + GetEventsFromStreams.new( + event_store: event_store, + stream_names: stream_names, + cursor: params["cursor"], + sort: sort, + ) + json( + html: + Renderer.new.render( + "swimlane/_rows", + urls: urls, + stream_names: stream_names, + events: reader.events, + sort: sort, + ), + more_url: (urls.swimlane_more_url(stream_names, reader.next_cursor, sort) if reader.more?), + ) + end + extensions.each do |extension| extension.register_routes(router, ExtensionContext.new(event_store, method(:extension_stylesheets), method(:extension_scripts))) end @@ -223,6 +258,16 @@ def html(body) [200, { "content-type" => "text/html;charset=utf-8" }, [body]] end + def json(body) + [200, { "content-type" => "application/json" }, [JSON.generate(body)]] + end + + def swimlane_params(params) + stream_names = Array(params["streams"]).reject { |name| name.nil? || name.empty? }.uniq + sort = ("as_of" if params["sort"] == "as_of") + [stream_names, sort] + end + def not_found(urls) renderer = Renderer.new content = renderer.render("not_found") diff --git a/ruby_event_store-browser/lib/ruby_event_store/browser/get_events_from_streams.rb b/ruby_event_store-browser/lib/ruby_event_store/browser/get_events_from_streams.rb new file mode 100644 index 0000000000..309fe21f0f --- /dev/null +++ b/ruby_event_store-browser/lib/ruby_event_store/browser/get_events_from_streams.rb @@ -0,0 +1,89 @@ +# frozen_string_literal: true + +module RubyEventStore + module Browser + GetEventsFromStreams = + Struct.new(:event_store, :stream_names, :cursor, :sort, :count, keyword_init: true) do + def initialize(event_store:, stream_names:, cursor: nil, sort: nil, count: PAGE_SIZE) + super(event_store: event_store, stream_names: stream_names, cursor: cursor, sort: sort, count: count) + end + + def events + @events ||= build_page + end + + def more? + events + @more + end + + def next_cursor + events.last && time_of(events.last.last).iso8601(TIMESTAMP_PRECISION) + end + + private + + def build_page + rows = stream_names.flat_map { |name| chunks[name].map { |event| [name, event] } } + @more = full_chunks.any? + return group(rows) unless @more + + boundary = full_chunks.map { |name| time_of(chunks[name].last) }.max + page_rows = rows.select { |_, event| time_of(event) >= boundary } + group(page_rows + drained_rows(boundary)) + end + + def chunks + @chunks ||= stream_names.to_h { |name| [name, read_chunk(name)] } + end + + def read_chunk(name) + scope = time_sorted(stream_scope(name)).backward.limit(count) + scope = scope.older_than(cursor_time) if cursor_time + scope.to_a + end + + def full_chunks + stream_names.select { |name| chunks[name].size == count } + end + + def drained_rows(boundary) + full_chunks + .select { |name| time_of(chunks[name].last) == boundary } + .flat_map { |name| drain(name, boundary).map { |event| [name, event] } } + end + + def drain(name, boundary) + time_sorted(stream_scope(name)).between(boundary..boundary).to_a + end + + def stream_scope(name) + name == SERIALIZED_GLOBAL_STREAM_NAME ? event_store.read : event_store.read.stream(name) + end + + def group(rows) + rows + .group_by { |_, event| event.event_id } + .map { |_, list| [list.map(&:first).uniq, list.first.last] } + .sort_by { |_, event| [time_of(event), event.event_id] } + .reverse + end + + def time_sorted(scope) + as_of? ? scope.as_of : scope.as_at + end + + def time_of(event) + event.metadata.fetch(as_of? ? :valid_at : :timestamp) + end + + def as_of? + sort == "as_of" + end + + def cursor_time + @cursor_time ||= cursor && Time.iso8601(cursor) + end + end + end +end diff --git a/ruby_event_store-browser/lib/ruby_event_store/browser/urls.rb b/ruby_event_store-browser/lib/ruby_event_store/browser/urls.rb index f0b778e884..eff3eb92df 100644 --- a/ruby_event_store-browser/lib/ruby_event_store/browser/urls.rb +++ b/ruby_event_store-browser/lib/ruby_event_store/browser/urls.rb @@ -39,6 +39,14 @@ def stream_page_url(stream_name, cursor, count) "#{stream_url(stream_name)}?#{query}" end + def swimlane_url(stream_names, sort = nil) + "#{app_url}/swimlane?#{swimlane_query(stream_names, sort)}" + end + + def swimlane_more_url(stream_names, cursor, sort) + "#{app_url}/swimlane/more?#{swimlane_query(stream_names, sort, [["cursor", cursor]])}" + end + def browser_js_url "#{app_url}/#{BROWSER_JS}" end @@ -50,6 +58,15 @@ def browser_css_url def ==(other) self.class.eql?(other.class) && app_url.eql?(other.app_url) end + + private + + def swimlane_query(stream_names, sort, extra = []) + pairs = stream_names.map { |name| ["streams[]", name] } + pairs.concat(extra) + pairs << ["sort", sort] if sort + URI.encode_www_form(pairs) + end end end end diff --git a/ruby_event_store-browser/lib/ruby_event_store/browser/views/streams/show.html.erb b/ruby_event_store-browser/lib/ruby_event_store/browser/views/streams/show.html.erb index ea9776f2e1..0d2aa5b600 100644 --- a/ruby_event_store-browser/lib/ruby_event_store/browser/views/streams/show.html.erb +++ b/ruby_event_store-browser/lib/ruby_event_store/browser/views/streams/show.html.erb @@ -15,9 +15,11 @@ - <% if extension_links.any? -%> + <% swimlane_link = { label: "Swimlane view", url: urls.swimlane_url([stream_name]) } unless stream_name == SERIALIZED_GLOBAL_STREAM_NAME -%> + <% links = [swimlane_link, *extension_links].compact -%> + <% if links.any? -%>
- <% extension_links.each do |link| -%> + <% links.each do |link| -%> <%= h(link[:label]) %> <% end -%>
diff --git a/ruby_event_store-browser/lib/ruby_event_store/browser/views/swimlane/_rows.html.erb b/ruby_event_store-browser/lib/ruby_event_store/browser/views/swimlane/_rows.html.erb new file mode 100644 index 0000000000..791c0aaaeb --- /dev/null +++ b/ruby_event_store-browser/lib/ruby_event_store/browser/views/swimlane/_rows.html.erb @@ -0,0 +1,20 @@ +<% time_field = sort == "as_of" ? :valid_at : :timestamp -%> +<% events.each_with_index do |(names, event), index| -%> + <% next_event = events[index + 1]&.last -%> + <% same_group = next_event && next_event.metadata[time_field].utc.iso8601(3) == event.metadata[time_field].utc.iso8601(3) -%> + <% cell = "align-middle px-2 py-1#{" border-b border-gray-300" unless same_group}" -%> + + + <% unless same_group -%> + <%= event.metadata[time_field].utc.strftime("%Y-%m-%dT%H:%M:%S.%3N") %> + <% end -%> + + <% stream_names.each_with_index do |stream_name, column| -%> + <%= " hover:bg-gray-100" if names.include?(stream_name) %>"> + <% if names.include?(stream_name) -%> + <%= h(event.event_type) %> + <% end -%> + + <% end -%> + +<% end -%> diff --git a/ruby_event_store-browser/lib/ruby_event_store/browser/views/swimlane/show.html.erb b/ruby_event_store-browser/lib/ruby_event_store/browser/views/swimlane/show.html.erb new file mode 100644 index 0000000000..d4b472f90e --- /dev/null +++ b/ruby_event_store-browser/lib/ruby_event_store/browser/views/swimlane/show.html.erb @@ -0,0 +1,50 @@ +
+
+

Comparing <%= h(stream_names.join(", ")) %>

+ +
+
+ Sort by + <% if sort == "as_of" -%> + created at + valid at + <% else -%> + created at + ">valid at + <% end -%> +
+ +
+ + + + + <% stream_names.each_with_index do |stream_name, column| -%> + + <% end -%> + + + + <%= render("swimlane/_rows", urls: urls, stream_names: stream_names, events: events, sort: sort) %> + +
Time"> +
+ <%= h(stream_name) %> + <% if stream_names.size > 1 -%> + × + <% end -%> +
+
+
+ + + +
+
+
+ +
+
+
+
+
diff --git a/ruby_event_store-browser/public/ruby_event_store_browser.js b/ruby_event_store-browser/public/ruby_event_store_browser.js index e68d377595..5899b5eb1a 100644 --- a/ruby_event_store-browser/public/ruby_event_store_browser.js +++ b/ruby_event_store-browser/public/ruby_event_store_browser.js @@ -92,3 +92,107 @@ application.register( } }, ) + +application.register( + "swimlane", + class extends Controller { + static targets = ["tbody", "time"] + static values = { moreUrl: String } + + get storageKey() { + return "ruby_event_store_browser.timezone" + } + + connect() { + this.onScroll = () => this.catchUp() + window.addEventListener("scroll", this.onScroll, { passive: true }) + this.onZoneChange = (event) => { + if (event.target.matches('[data-timezone-target="select"]')) this.renderTimes(event.target.value) + } + document.addEventListener("change", this.onZoneChange) + this.renderTimes(this.zone()) + this.catchUp() + } + + disconnect() { + window.removeEventListener("scroll", this.onScroll) + document.removeEventListener("change", this.onZoneChange) + } + + catchUp() { + if (!this.moreUrlValue) return + if (document.body.scrollHeight - (window.scrollY + window.innerHeight) > 200) return + this.loadMore() + } + + loadMore() { + const url = this.moreUrlValue + if (!url) return + this.moreUrlValue = "" + + fetch(url, { headers: { Accept: "application/json" } }) + .then((response) => response.json()) + .then(({ html, more_url }) => { + this.tbodyTarget.insertAdjacentHTML("beforeend", html) + this.moreUrlValue = more_url || "" + this.renderTimes(this.zone()) + this.catchUp() + }) + } + + zone() { + const stored = localStorage.getItem(this.storageKey) + const detected = Intl.DateTimeFormat().resolvedOptions().timeZone + try { + Intl.DateTimeFormat("en-US", { timeZone: stored }) + return stored || detected + } catch (_) { + return detected + } + } + + renderTimes(tz) { + this.timeTargets.forEach((el) => { + el.textContent = this.format(el.dataset.iso, tz) + el.setAttribute("title", tz) + }) + } + + format(iso, tz) { + const parts = Object.fromEntries( + new Intl.DateTimeFormat("en-US", { + timeZone: tz, year: "numeric", month: "2-digit", day: "2-digit", + hour: "2-digit", minute: "2-digit", second: "2-digit", + fractionalSecondDigits: 3, hour12: false, + }).formatToParts(new Date(iso)).map((p) => [p.type, p.value]) + ) + return `${parts.year}-${parts.month}-${parts.day}T${parts.hour}:${parts.minute}:${parts.second}.${parts.fractionalSecond}` + } + }, +) + +application.register( + "swimlane-add", + class extends Controller { + static targets = ["dialog", "input"] + + open(event) { + event?.preventDefault() + this.dialogTarget.showModal() + this.inputTarget.focus() + } + + close() { + this.dialogTarget.close() + } + + go(event) { + event.preventDefault() + const name = this.inputTarget.value + if (!name) return + const url = new URL(window.location.href) + url.searchParams.append("streams[]", name) + window.location = url.toString() + } + }, +) diff --git a/ruby_event_store-browser/spec/get_events_from_streams_spec.rb b/ruby_event_store-browser/spec/get_events_from_streams_spec.rb new file mode 100644 index 0000000000..45d911f070 --- /dev/null +++ b/ruby_event_store-browser/spec/get_events_from_streams_spec.rb @@ -0,0 +1,218 @@ +# frozen_string_literal: true + +require "spec_helper" + +module RubyEventStore + module Browser + ::RSpec.describe GetEventsFromStreams do + let(:event_store) { RubyEventStore::Client.new } + let(:base_time) { Time.utc(2024, 1, 1, 12, 0, 0) } + + def reader(stream_names:, cursor: nil, sort: nil, count: 20) + GetEventsFromStreams.new( + event_store: event_store, + stream_names: stream_names, + cursor: cursor, + sort: sort, + count: count, + ) + end + + def next_page(previous, stream_names:, count:, sort: nil) + reader(stream_names: stream_names, cursor: previous.next_cursor, sort: sort, count: count) + end + + def publish(stream:, at: nil, valid_at: nil) + event = DummyEvent.new + if at + event.metadata[:timestamp] = at + event.metadata[:valid_at] = valid_at || at + end + event_store.publish(event, stream_name: stream) + event + end + + specify "merges events from multiple streams, newest first, tagged with their owning stream" do + a1 = publish(stream: "a", at: base_time + 1) + b1 = publish(stream: "b", at: base_time + 2) + a2 = publish(stream: "a", at: base_time + 3) + + result = reader(stream_names: %w[a b]).events.map { |names, event| [names, event.event_id] } + expect(result).to eq([[["a"], a2.event_id], [["b"], b1.event_id], [["a"], a1.event_id]]) + end + + specify "orders by event time, not by position in the stream" do + newer = publish(stream: "a", at: base_time + 2) + older = publish(stream: "a", at: base_time + 1) + + result = reader(stream_names: %w[a]).events.map { |_, event| event.event_id } + expect(result).to eq([newer.event_id, older.event_id]) + end + + specify "a single stream page is capped at count" do + 6.times { |i| publish(stream: "a", at: base_time + i) } + expect(reader(stream_names: %w[a], count: 3).events.size).to eq(3) + end + + specify "a page emits everything above the completeness horizon, not just count events" do + a1 = publish(stream: "a", at: base_time + 1) + b2 = publish(stream: "b", at: base_time + 2) + a3 = publish(stream: "a", at: base_time + 3) + b4 = publish(stream: "b", at: base_time + 4) + + page = reader(stream_names: %w[a b], count: 2) + expect(page.events.map { |_, event| event.event_id }).to eq([b4, a3, b2].map(&:event_id)) + expect(page.more?).to eq(true) + + rest = next_page(page, stream_names: %w[a b], count: 2) + expect(rest.events.map { |_, event| event.event_id }).to eq([a1.event_id]) + expect(rest.more?).to eq(false) + end + + specify "more? is true when the page is full" do + 3.times { |i| publish(stream: "a", at: base_time + i) } + expect(reader(stream_names: %w[a], count: 3).more?).to eq(true) + end + + specify "streams read to exhaustion merge into a single page" do + 2.times { |i| publish(stream: "a", at: base_time + i) } + 2.times { |i| publish(stream: "b", at: base_time + 10 + i) } + + page = reader(stream_names: %w[a b], count: 3) + expect(page.events.size).to eq(4) + expect(page.more?).to eq(false) + end + + specify "more? is false once every stream is exhausted" do + 3.times { |i| publish(stream: "a", at: base_time + i) } + expect(reader(stream_names: %w[a], count: 20).more?).to eq(false) + end + + specify "next_cursor is the last returned event's timestamp" do + publish(stream: "a", at: base_time + 1) + a2 = publish(stream: "a", at: base_time + 2) + a3 = publish(stream: "a", at: base_time + 3) + + page = reader(stream_names: %w[a], count: 2) + expect(page.events.map { |_, event| event.event_id }).to eq([a3.event_id, a2.event_id]) + expect(page.next_cursor).to eq((base_time + 2).iso8601(TIMESTAMP_PRECISION)) + end + + specify "continues below the cursor timestamp" do + a1 = publish(stream: "a", at: base_time + 1) + a2 = publish(stream: "a", at: base_time + 2) + a3 = publish(stream: "a", at: base_time + 3) + + page = reader(stream_names: %w[a], count: 2) + rest = next_page(page, stream_names: %w[a], count: 2) + expect(page.events.map { |_, event| event.event_id }).to eq([a3.event_id, a2.event_id]) + expect(rest.events.map { |_, event| event.event_id }).to eq([a1.event_id]) + expect(rest.more?).to eq(false) + end + + specify "does not skip an event from another stream sharing the boundary timestamp" do + old = publish(stream: "a", at: base_time + 1) + tied_a = publish(stream: "a", at: base_time + 2) + tied_b = publish(stream: "b", at: base_time + 2) + newest = publish(stream: "b", at: base_time + 3) + + page = reader(stream_names: %w[a b], count: 2) + rest = next_page(page, stream_names: %w[a b], count: 2) + + expect(page.events.map { |_, event| event.event_id }).to match_array( + [newest, tied_a, tied_b].map(&:event_id), + ) + expect(rest.events.map { |_, event| event.event_id }).to eq([old.event_id]) + end + + specify "a page runs to the end of its last timestamp instead of splitting the group" do + published = 7.times.map { publish(stream: "a", at: base_time) } + + page = reader(stream_names: %w[a], count: 3) + expect(page.events.map { |_, event| event.event_id }).to match_array(published.map(&:event_id)) + + rest = next_page(page, stream_names: %w[a], count: 3) + expect(rest.events).to eq([]) + expect(rest.more?).to eq(false) + end + + specify "pages through mixed unique and tied timestamps without drops or duplicates" do + published = 3.times.map { |i| publish(stream: "a", at: base_time + i) } + published += 4.times.map { publish(stream: "a", at: base_time + 10) } + published += 2.times.map { |i| publish(stream: "b", at: base_time + 20 + i) } + + pages = [reader(stream_names: %w[a b], count: 3)] + pages << next_page(pages.last, stream_names: %w[a b], count: 3) while pages.last.more? + + emitted = pages.flat_map { |page| page.events.map { |_, event| event.event_id } } + expect(emitted).to match_array(published.map(&:event_id)) + end + + specify "an event linked into more than one compared stream is one row shown in both columns, not duplicated" do + a1 = publish(stream: "a", at: base_time) + event_store.link([a1.event_id], stream_name: "b") + + result = reader(stream_names: %w[a b]).events + expect(result.map { |_, event| event.event_id }).to eq([a1.event_id]) + expect(result.first.first).to match_array(%w[a b]) + end + + specify "does not repeat a stream in the columns of an event fetched twice" do + 3.times { publish(stream: "a", at: base_time) } + + result = reader(stream_names: %w[a], count: 3).events + expect(result.map(&:first)).to eq([["a"], ["a"], ["a"]]) + end + + specify "sorts by validity time when sort is as_of" do + e1 = publish(stream: "a", at: base_time + 3, valid_at: base_time + 11) + e2 = publish(stream: "b", at: base_time + 2, valid_at: base_time + 12) + e3 = publish(stream: "a", at: base_time + 1, valid_at: base_time + 13) + + by_append = reader(stream_names: %w[a b]).events.map { |_, event| event.event_id } + by_validity = reader(stream_names: %w[a b], sort: "as_of").events.map { |_, event| event.event_id } + expect(by_append).to eq([e1.event_id, e2.event_id, e3.event_id]) + expect(by_validity).to eq([e3.event_id, e2.event_id, e1.event_id]) + end + + specify "as_of cursor pages along the validity axis" do + e1 = publish(stream: "a", at: base_time + 3, valid_at: base_time + 11) + e2 = publish(stream: "a", at: base_time + 2, valid_at: base_time + 12) + e3 = publish(stream: "a", at: base_time + 1, valid_at: base_time + 13) + + page = reader(stream_names: %w[a], sort: "as_of", count: 2) + expect(page.events.map { |_, event| event.event_id }).to eq([e3.event_id, e2.event_id]) + expect(page.next_cursor).to eq((base_time + 12).iso8601(TIMESTAMP_PRECISION)) + + rest = next_page(page, stream_names: %w[a], sort: "as_of", count: 2) + expect(rest.events.map { |_, event| event.event_id }).to eq([e1.event_id]) + expect(rest.more?).to eq(false) + end + + specify "as_of pages run to the end of their last validity time instead of splitting the group" do + published = 4.times.map { |i| publish(stream: "a", at: base_time + i, valid_at: base_time + 10) } + + page = reader(stream_names: %w[a], sort: "as_of", count: 3) + expect(page.events.map { |_, event| event.event_id }).to match_array(published.map(&:event_id)) + end + + specify "the global stream alias reads all events, not a stream named all" do + a1 = publish(stream: "orders", at: base_time + 1) + a2 = publish(stream: "shipments", at: base_time + 2) + + result = reader(stream_names: %w[all orders]).events.map { |names, event| [names, event.event_id] } + expect(result).to eq([[%w[all], a2.event_id], [%w[all orders], a1.event_id]]) + end + + specify "reads each stream once per page" do + publish(stream: "a", at: base_time) + publish(stream: "b", at: base_time + 1) + + allow(event_store).to receive(:read).and_call_original + reader(stream_names: %w[a b c]).events + + expect(event_store).to have_received(:read).exactly(3).times + end + end + end +end diff --git a/ruby_event_store-browser/spec/swimlane_spec.rb b/ruby_event_store-browser/spec/swimlane_spec.rb new file mode 100644 index 0000000000..0bccfe9a7d --- /dev/null +++ b/ruby_event_store-browser/spec/swimlane_spec.rb @@ -0,0 +1,193 @@ +# frozen_string_literal: true + +require "spec_helper" +require "json" +require "uri" + +module RubyEventStore + module Browser + ::RSpec.describe Browser do + let(:event_store) { RubyEventStore::Client.new } + let(:app) { App.for(event_store_locator: -> { event_store }) } + let(:client) { Rack::MockRequest.new(Rack::Lint.new(app)) } + let(:base_time) { Time.utc(2024, 1, 1, 12, 0, 0) } + + def event_at(time, valid_at: nil) + DummyEvent.new.tap do |event| + event.metadata[:timestamp] = time + event.metadata[:valid_at] = valid_at if valid_at + end + end + + def more_url(stream_names, cursor_time, sort = nil) + pairs = stream_names.map { |name| ["streams[]", name] } + pairs << ["cursor", cursor_time.iso8601(TIMESTAMP_PRECISION)] + pairs << ["sort", sort] if sort + "http://example.org/swimlane/more?#{URI.encode_www_form(pairs)}" + end + + specify "compare view lists events from all compared streams" do + event_store.append(e1 = event_at(base_time + 1), stream_name: "fizz") + event_store.append(e2 = event_at(base_time + 2), stream_name: "buzz") + + response = client.get("/swimlane?streams%5B%5D=fizz&streams%5B%5D=buzz") + expect(response.status).to eq(200) + expect(response.headers["content-type"]).to eq("text/html;charset=utf-8") + expect(response.body).to include("Comparing fizz, buzz") + expect(response.body).to include(e1.event_id) + expect(response.body).to include(e2.event_id) + expect(response.body.scan("") + end + + specify "compare view exposes the url for fetching older events" do + event_store.append(event_at(base_time), stream_name: "buzz") + events = Array.new(Browser::PAGE_SIZE + 1) { |i| event_at(base_time + i + 1) } + event_store.append(events, stream_name: "fizz") + + body = client.get("/swimlane?streams%5B%5D=fizz&streams%5B%5D=buzz").body + expect(body).to include( + "data-swimlane-more-url-value=\"#{more_url(%w[fizz buzz], base_time + 2)}\"", + ) + end + + specify "compare view has no more-url when everything fits one page" do + event_store.append(event_at(base_time), stream_name: "fizz") + + expect(client.get("/swimlane?streams%5B%5D=fizz").body).to include(%q[data-swimlane-more-url-value=""]) + end + + specify "compare more endpoint returns next rows and cursor as JSON" do + event_store.append(event_at(base_time), stream_name: "buzz") + events = Array.new(Browser::PAGE_SIZE + 1) { |i| event_at(base_time + i + 1) } + event_store.append(events, stream_name: "fizz") + + response = client.get("/swimlane/more?streams%5B%5D=fizz&streams%5B%5D=buzz") + expect(response.status).to eq(200) + expect(response.headers["content-type"]).to eq("application/json") + payload = JSON.parse(response.body) + expect(payload.keys).to eq(%w[html more_url]) + expect(payload["html"]).to include(events.last.event_id) + expect(payload["html"]).not_to include(" "", "more_url" => nil }) + end + + specify "compare more endpoint pages from the cursor" do + e1, e2, e3 = Array.new(3) { |i| event_at(base_time + i + 1) } + event_store.append([e1, e2, e3], stream_name: "fizz") + + query = URI.encode_www_form([["streams[]", "fizz"], ["cursor", (base_time + 2).iso8601(TIMESTAMP_PRECISION)]]) + payload = JSON.parse(client.get("/swimlane/more?#{query}").body) + expect(payload["html"]).to include(e1.event_id) + expect(payload["html"]).not_to include(e2.event_id) + expect(payload["html"]).not_to include(e3.event_id) + expect(payload["more_url"]).to be_nil + end + + specify "compare view sorts by validity time when requested" do + event_store.append(e1 = event_at(base_time + 2, valid_at: base_time + 11), stream_name: "fizz") + event_store.append(e2 = event_at(base_time + 1, valid_at: base_time + 12), stream_name: "buzz") + + by_append = client.get("/swimlane?streams%5B%5D=fizz&streams%5B%5D=buzz").body + expect(by_append.index(e1.event_id)).to be < by_append.index(e2.event_id) + expect(by_append).to include("/swimlane?streams%5B%5D=fizz&streams%5B%5D=buzz&sort=as_of").or include( + "/swimlane?streams%5B%5D=fizz&streams%5B%5D=buzz&sort=as_of", + ) + + by_validity = client.get("/swimlane?streams%5B%5D=fizz&streams%5B%5D=buzz&sort=as_of").body + expect(by_validity.index(e2.event_id)).to be < by_validity.index(e1.event_id) + expect(by_validity).to include((base_time + 12).utc.iso8601(6)) + end + + specify "compare more endpoint honors as_of sort and carries it in more_url" do + events = Array.new(Browser::PAGE_SIZE + 1) { |i| event_at(base_time + i, valid_at: base_time + 100 - i) } + event_store.append(events, stream_name: "fizz") + + query = URI.encode_www_form([["streams[]", "fizz"], ["sort", "as_of"]]) + payload = JSON.parse(client.get("/swimlane/more?#{query}").body) + expect(payload["html"]).to include(events[0].event_id) + expect(payload["html"]).to include((base_time + 100).utc.iso8601(6)) + expect(payload["html"]).not_to include(events[20].event_id) + expect(payload["more_url"]).to eq(more_url(%w[fizz], base_time + 81, "as_of")) + end + + specify "unknown sort values fall back to created-at order" do + event_store.append(e1 = event_at(base_time + 2, valid_at: base_time + 11), stream_name: "fizz") + event_store.append(e2 = event_at(base_time + 1, valid_at: base_time + 12), stream_name: "buzz") + + body = client.get("/swimlane?streams%5B%5D=fizz&streams%5B%5D=buzz&sort=wrong").body + expect(body.index(e1.event_id)).to be < body.index(e2.event_id) + + query = URI.encode_www_form([["streams[]", "fizz"], ["streams[]", "buzz"], ["sort", "wrong"]]) + html = JSON.parse(client.get("/swimlane/more?#{query}").body).fetch("html") + expect(html.index(e1.event_id)).to be < html.index(e2.event_id) + end + + specify "compare endpoints tolerate valueless stream params" do + event_store.append(event_at(base_time), stream_name: "fizz") + + expect(client.get("/swimlane?streams%5B%5D").body).to include("Comparing ") + expect(JSON.parse(client.get("/swimlane/more?streams%5B%5D").body)).to eq({ "html" => "", "more_url" => nil }) + end + + specify "compare view carries as_of sort in the url for fetching older events" do + events = Array.new(Browser::PAGE_SIZE + 1) { |i| event_at(base_time + i, valid_at: base_time + 100 - i) } + event_store.append(events, stream_name: "fizz") + + body = client.get("/swimlane?streams%5B%5D=fizz&sort=as_of").body + expect(body).to include("data-swimlane-more-url-value=\"#{more_url(%w[fizz], base_time + 81, "as_of")}\"") + end + + specify "comparing the all stream shows the global timeline" do + event_store.append(event = event_at(base_time), stream_name: "orders") + + body = client.get("/swimlane?streams%5B%5D=all").body + expect(body).to include(event.event_id) + end + + specify "stream page links to comparing that stream" do + event_store.append(event_at(base_time), stream_name: "fizz") + + body = client.get("/streams/fizz").body + expect(body).to include("Swimlane view") + expect(body).to include("http://example.org/swimlane?streams%5B%5D=fizz") + end + + specify "the global stream page does not link to the swimlane view" do + event_store.append(event_at(base_time), stream_name: "fizz") + + expect(client.get("/streams/all").body).not_to include("Swimlane view") + end + + specify "removing a stream from comparison keeps the sort" do + event_store.append(event_at(base_time), stream_name: "fizz") + event_store.append(event_at(base_time + 1), stream_name: "buzz") + + body = client.get("/swimlane?streams%5B%5D=fizz&streams%5B%5D=buzz&sort=as_of").body + expect(body).to include("/swimlane?streams%5B%5D=buzz&sort=as_of").or include( + "/swimlane?streams%5B%5D=buzz&sort=as_of", + ) + end + end + end +end diff --git a/ruby_event_store-browser/spec/urls_spec.rb b/ruby_event_store-browser/spec/urls_spec.rb index fe96eb675b..0d15e7a9e2 100644 --- a/ruby_event_store-browser/spec/urls_spec.rb +++ b/ruby_event_store-browser/spec/urls_spec.rb @@ -90,6 +90,34 @@ module Browser ) end + specify "swimlane_url" do + urls = Urls.from_configuration("http://example.com:9393", "/res") + expect(urls.swimlane_url(%w[fizz buzz])).to eq( + "http://example.com:9393/res/swimlane?streams%5B%5D=fizz&streams%5B%5D=buzz", + ) + end + + specify "swimlane_url with sort" do + urls = Urls.from_configuration("http://example.com:9393", nil) + expect(urls.swimlane_url(%w[fizz], "as_of")).to eq( + "http://example.com:9393/swimlane?streams%5B%5D=fizz&sort=as_of", + ) + end + + specify "swimlane_more_url" do + urls = Urls.from_configuration("http://example.com:9393", nil) + expect(urls.swimlane_more_url(%w[fizz buzz], "2024-01-01T12:00:00.000000Z", nil)).to eq( + "http://example.com:9393/swimlane/more?streams%5B%5D=fizz&streams%5B%5D=buzz&cursor=2024-01-01T12%3A00%3A00.000000Z", + ) + end + + specify "swimlane_more_url with sort" do + urls = Urls.from_configuration("http://example.com:9393", nil) + expect(urls.swimlane_more_url(%w[fizz], "2024-01-01T12:00:00.000000Z", "as_of")).to eq( + "http://example.com:9393/swimlane/more?streams%5B%5D=fizz&cursor=2024-01-01T12%3A00%3A00.000000Z&sort=as_of", + ) + end + specify "browser_css_url" do urls = Urls.from_configuration("http://example.com:9393", "/res") expect(urls.browser_css_url).to eq("http://example.com:9393/res/#{BROWSER_CSS}")