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? -%>
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 -%>
+
+
+
+
+
+
+ | Time |
+ <% stream_names.each_with_index do |stream_name, column| -%>
+ ">
+
+ <%= h(stream_name) %>
+ <% if stream_names.size > 1 -%>
+ ×
+ <% end -%>
+
+ |
+ <% end -%>
+
+
+
+ <%= render("swimlane/_rows", urls: urls, stream_names: stream_names, events: events, sort: sort) %>
+
+
+
+
+
+
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}")