Correlation chain
Commands maintain a causation/correlation chain for traceability:
event = source_cmd.correlate(SomeEvent.new(payload: { ... }))
event.causation_id # => source_cmd.id
event.correlation_id # => source_cmd.correlation_id
Command-oriented web framework for Ruby. Server-rendered, async, reactive, multi-player by default.
A Ruby gem for building server-driven, reactive web applications. Sidereal combines a Rack-compatible router with an event-driven architecture using typed messages, commands, pages, SSE, and pub/sub.
Only one way to reason about business logic handling and UI updates, whether one-off or long-running, re-tryable tasks.
I talked about the motivation and techniques here.
Built on Datastar (SSE streaming, HTML morphing), Phlex (HTML rendering), Plumb (typed data), Async (fiber concurrency). Designed to run on Falcon.
See the examples directory for demos.
Add to your Gemfile:
bundle add sidereal
Or install directly:
gem install sidereal
A Sidereal app has three main parts: commands (typed data), command handlers (state changes), and pages (reactive UI).
require 'sidereal'
# 1. Define command messages
AddTodo = Sidereal::Message.define('todos.add') do
attribute :title, Sidereal::Types::String.present
end
# 2. Define a page
class TodoPage < Sidereal::Page
path '/'
# React to events by pushing HTML updates via SSE
on AddTodo do |evt|
browser.patch_elements load(params)
end
def self.load(_params, _ctx)
new(todos: TODOS.values)
end
def initialize(todos: [])
@todos = todos
end
def view_template
div do
# An Ajax form to dispatch a command to the backend
command AddTodo do |f|
f.text_field :title, placeholder: 'What needs to be done?'
button(type: :submit) { 'Add' }
end
ul do
@todos.each { |t| li { t.title } }
end
end
end
end
# 3. Wire it up in an App
TODOS = {}
class TodoApp < Sidereal::App
session secret: ENV.fetch('SESSION_SECRET')
# Expose AddTodo to the browser's POST /commands
handle AddTodo
# Register the async worker handler
command AddTodo do |cmd|
TODOS[cmd.id] = cmd.payload
end
page TodoPage
end
Commands are typed, immutable data objects defined with Message.define. Each message has an auto-generated UUID, a type string, timestamps, metadata, and a typed payload.
AddTodo = Sidereal::Message.define('todos.add') do
attribute :todo_id, Sidereal::Types::AutoUUID
attribute :title, Sidereal::Types::String.present
end
RemoveTodo = Sidereal::Message.define('todos.remove') do
attribute :todo_id, Sidereal::Types::UUID::V4
end
Notify = Sidereal::Message.define('todos.notify') do
attribute :message, String
end
Commands use dot-separated type strings (e.g. 'todos.add') for registry lookup and serialization. Payload attributes are validated using Plumb types.
cmd = AddTodo.new(payload: { title: 'Buy milk' })
cmd.id # => "a1b2c3d4-..."
cmd.type # => "todos.add"
cmd.payload.title # => "Buy milk"
cmd.created_at # => 2026-03-26 10:00:00 +0000
Commands maintain a causation/correlation chain for traceability:
event = source_cmd.correlate(SomeEvent.new(payload: { ... }))
event.causation_id # => source_cmd.id
event.correlation_id # => source_cmd.correlation_id
Sidereal::App is a web router with that implements a full reactive framework: command handlers, pages, layouts, and SSE streaming. It automatically sets up POST /commands and GET /updates/:channel_name endpoints.
class ChatApp < Sidereal::App
session secret: ENV.fetch('SESSION_SECRET')
layout ChatLayout
# SendMessage is submitted from the browser
# Use .handle to white-list this command in the HTTP endpoint
handle SendMessage
# Incoming command is dispatched to a background fiber,
# and handled by this block
# The block can optional #dispatch a new command in a workflow
command SendMessage do |cmd|
MessageLog.append(cmd)
dispatch Notify, message: "#{cmd.payload.author}: #{cmd.payload.content}"
end
# Notify is dispatched internally and never exposed to the web
command Notify do |cmd|
# no-op, but events from this command will still be published
end
# Mount a page to be served on ChatPage.path
page ChatPage
end
Commands are split into two registrations:
command registers an async handler with the app’s Commander. Worker fibers pick the command off the store and run this block. Commands registered only via command are internal — they can be produced by other handlers, automations, or sagas, but cannot be submitted from the browser.handle exposes a command to POST /commands (see Custom command handlers). Any type that isn’t handle-registered returns 404 on POST.Inside a command block, use dispatch to produce events or enqueue follow-up commands.
# Internal command — dispatched from other handlers, never from the browser
SendEmail = Sidereal::Message.define('mail.send') { attribute :to, String }
command SendEmail do |cmd|
Mailer.deliver(cmd.payload.to)
end
# Web-facing command — exposed via `handle`, processed async via `command`
handle AddTodo
command AddTodo do |cmd|
TODOS[cmd.payload.todo_id] = cmd.payload
# Dispatching a registered command type enqueues it for processing
dispatch SendEmail, to: 'user@example.com'
# Dispatching any other message type produces a transient event
dispatch Notify, message: "Added: #{cmd.payload.title}"
end
Use broadcast inside a command handler to publish a message immediately to the SSE stream, without waiting for the command to finish processing. Useful for progress indicators.
command AskLLM do |cmd|
broadcast Working # immediately tells the UI "thinking..."
response = llm.ask(cmd.payload.content)
dispatch SendMessage, author: 'Bot', content: response.content
end
Define helper methods available inside command handlers:
class ChatApp < Sidereal::App
command_helpers do
private def chat
@chat ||= RubyLLM.chat
end
end
command AskLLM do |cmd|
response = chat.ask(cmd.payload.content)
dispatch SendMessage, author: 'Bot', content: response.content
end
end
handle declares which commands the browser is allowed to submit to POST /commands. Types not registered with handle return 404.
handle accepts one or more command classes. Called without a block, it installs the default handler: validate the command, append it to the async store, and return 200. Worker fibers then pick it up and run the matching command block.
class TodoApp < Sidereal::App
# Expose multiple commands at once with the default handler
handle AddTodo, RemoveTodo
command AddTodo do |cmd| # async worker handler
TODOS[cmd.payload.todo_id] = cmd.payload
end
end
Pass a block to handle to process the command synchronously during the HTTP request instead — useful for lightweight mutations, or when you want to stream DOM updates back to the browser immediately:
handle AddTodo do |cmd|
TODOS[cmd.id] = cmd.payload.to_h
browser.patch_elements TodoList.new(TODOS.values)
end
A custom handle block replaces the default async-dispatch behaviour. If you still want the async worker to run, call dispatch(cmd) inside the block.
handle AddTodo do |cmd|
browser.patch_elements %(<p id="notification">Processing...</p>)
dispatch(cmd) # <= schedule command for background processing
end
handle does not register the command with the async Commander. To have a web-submitted command also processed by workers, pair handle with a command block as shown above.
Inside a handle block you have access to:
browser — the SSE stream for pushing DOM updates (see SSE reactions for the full API)dispatch(MessageClass, payload) — correlate and append a follow-up command to the async storepatch_command_errors(errors) — stream field-level validation errors back to the formstore, pubsub, params, session — the usual App instance helpersWhen the browser submits a command via Datastar (the default), the request accepts SSE responses. The handle block can use browser to push HTML patches, signal updates, or JavaScript execution — just like page on reactions:
handle AddTodo do |cmd|
TODOS[cmd.id] = cmd.payload.to_h
browser.patch_elements TodoList.new(TODOS.values)
end
browser is an alias for the Datastar dispatcher — each call above produces a one-off SSE event. For multi-step, real-time updates over a single request, use browser.stream { |sse| ... }:
handle GenerateReport do |cmd|
browser.stream do |sse|
sse.patch_elements %(<div id="status">Working...</div>)
expensive_work do |progress|
sse.patch_signals progress: progress
end
sse.patch_elements %(<div id="status">Done</div>)
end
end
The handle block runs synchronously before any SSE streaming starts, so session[:x] = … writes inside it commit to the session cookie as expected.
Use dispatch to enqueue a command for async processing. The dispatched command is automatically correlated to the source command:
handle AddTodo do |cmd|
TODOS[cmd.id] = cmd.payload.to_h
browser.patch_elements TodoList.new(TODOS.values)
dispatch NotifyUser, text: "Todo added: #{cmd.payload.title}"
end
Use patch_command_errors to stream field-level errors back to the form. This works with the command form component, which generates the matching element IDs:
handle PlaceOrder do |cmd|
errors = validate_stock(cmd.payload)
if errors.any?
patch_command_errors(errors)
else
ORDERS[cmd.id] = cmd.payload.to_h
browser.patch_elements OrderConfirmation.new(cmd.payload)
end
end
Some things need more than one macro to wire in: a command to expose with handle, a commander to register with commands, maybe a channel resolver or a page. install lets the extensions do that itself, in one line, so the knowledge of what it needs stays with it:
class DataflowApp < Sidereal::App
install Onboarding, reset: true
end
App.install(installer, ...) calls installer.sidereal_install(app, ...), passing the app class and any extra arguments through. Anything that responds to sidereal_install qualifies; inside it the installer uses the app’s own macros:
module Billing
def self.sidereal_install(app, webhooks: true)
app.handle(PlaceOrder)
app.commands(Billing::Commander)
app.channel_name(OrderPlaced) { |evt| "orders.#{evt.payload.order_id}" }
app.page(WebhookLogPage) if webhooks
end
end
class ShopApp < Sidereal::App
install Billing, webhooks: false
end
install raises ArgumentError for an object without sidereal_install, and returns the app so it chains like the other macros.
The component helper renders any object responding to #call(context:), passing the app instance as context. This is how Sidereal pages and layouts are rendered under the hood, and it’s available inside any route block defined on an App subclass.
class MyApp < Sidereal::App
get '/dashboard' do
component DashboardPage.new(current_user)
end
get '/error' do
component ErrorPage.new, status: 422
end
end
Pages are reactive Phlex components that re-render parts of the UI in response to events via SSE.
class TodoPage < Sidereal::Page
path '/'
# React to events -- re-render components via SSE
on AddTodo do |evt|
browser.patch_elements TodoList.new(TODOS.values)
end
on RemoveTodo do |evt|
browser.patch_elements TodoList.new(TODOS.values)
end
on Notify do |evt|
browser.patch_elements ActivityItem.new(evt), mode: 'append', selector: '#feed'
end
# Load is called on initial page render and on SSE reconnect
def self.load(_params, _ctx)
new(todos: TODOS.values)
end
def initialize(todos: [])
@todos = todos
end
def view_template
div do
render TodoList.new(@todos)
aside do
h2 { 'Activity' }
div(id: 'feed')
end
end
end
end
Inside an on block, you have access to:
browser – the SSE stream for pushing updatesload(params) – re-instantiate the page with current dataparams – the current page params from Datastar signalson AddTodo do |evt|
# Replace an element's content with a re-rendered component
browser.patch_elements load(params)
# Or target a specific element
browser.patch_elements TodoList.new(TODOS.values)
# Append to a container
browser.patch_elements ActivityItem.new(evt), mode: 'append', selector: '#feed'
# Patch signal values
browser.patch_signals progress: 99
# Execute JavaScript on the client
browser.execute_script %(scrollToBottom('messages'))
end
See more about this here.
A page that renders a form for a command almost always wants to re-render when that command comes back over pubsub. You don’t have to write that reaction: every command form rendered inside a page, at any depth of its component tree, registers on CommandClass on that page. on without a block means “reload the page”, so the two pages below react the same way:
class TodoPage < Sidereal::Page
on AddTodo # reload when AddTodo comes back
def view_template
command AddTodo do |f|
f.text_field :title
end
end
end
class TodoPage < Sidereal::Page
def view_template
command AddTodo do |f| # registers `on AddTodo` on TodoPage as it renders
f.text_field :title
end
end
end
This is the default. A page that would rather say exactly what it reacts to opts out with disable_causal_reactivity!, which stops its forms from registering anything, and keeps every on it declares, block or not:
class AuditPage < Sidereal::Page
disable_causal_reactivity!
on AuditEntryAdded # only this, however many forms the page renders
end
Subclasses inherit the setting.
The match is on the message’s own type or on its correlation_type. Every message carries the type of the message at the root of its causal chain (Sourced::Message#correlation_type, recorded in metadata[:correlation_type] by #correlate), so a page registered for AddTodo also reloads on the events a handler produced from it, the commands those events triggered, and so on. This is what makes the same page work over a backend that publishes the command itself and one, like Sourced, that publishes only the resulting events. Naming an event or a projector signal works too: on GamesProjector::Projected reloads on every signal that projector publishes, whichever command’s chain it belongs to.
on also takes anything that answers sidereal_events with a list of message classes, and registers each of them. A Sidereal::Commander answers with the commands it handles, which the dispatcher publishes along with everything dispatched from them, so on MyCommander covers all of it. With the Sourced integration loaded, a decider answers with the events it evolves, its own and any foreign ones that change its state, and a projector with its Projected signal. So an event-sourced page names the reactor it follows:
class DonationPage < Sidereal::Page
on Donation # every event that changes a donation
end
class HomePage < Sidereal::Page
on GamesProjector # the lobby's read model committed
end
A handler written with a block always wins for messages of exactly its class. The page checks reactions first and only falls back to the reload when the message’s class has no handler of its own, so on TodoAdded do |evt| ... end next to a rendered command AddTodo runs the block for TodoAdded and reloads for anything else in that chain. At most one of the two runs per message.
Each page subscribes to a single PubSub channel via GET /updates/:channel_name. The default is 'system', which means every page receives every published event. Override Page#channel_name to scope a page’s SSE stream to a narrower topic — for example, “only events for this donation” or “only events for this chat room”.
class DonationPage < Sidereal::Page
path '/:donation_id'
def initialize(donation_id:, **)
@donation_id = donation_id
end
# Each donation page only receives events on its own channel
def channel_name = "donations.#{@donation_id}"
end
For events to actually reach that channel, declare how the App derives a channel name from each message. The channel_name macro registers a resolver on the process-global Sidereal.channels registry. Pass message classes positionally to scope the resolver, or no arguments to register a catch-all:
class DonationsApp < Sidereal::App
# Catch-all: every message goes through this block
channel_name do |msg|
"donations.#{msg.payload.donation_id}"
end
# Or scope to specific message classes:
channel_name SelectAmount, EnterDonorDetails do |msg|
"donations.#{msg.payload.donation_id}"
end
handle SelectAmount
end
The block runs for every message the dispatcher publishes — both the incoming command and the events it emits. Resolution is O(1) per message: typed registrations win first, then the catch-all, then a fallback to the literal 'system' channel (so an app that registers nothing still publishes successfully). System notifications (Sidereal::System::NotifyRetry/NotifyFailure) are pre-routed via the :source_channel metadata that the dispatcher stamps; user-supplied resolvers never see them.
Channel routing also works outside the App class — call Sidereal.channels.channel_name(...) from anywhere (e.g. a dedicated routes file) for apps where the registrations grow large enough to warrant their own home.
The registry locks itself once boot is over: Sidereal::Falcon::Environment::Service calls Sidereal.channels.lock! after class-loading and pubsub startup, before workers start consuming. Subsequent channel_name(...) calls raise Sidereal::Channels::LockedError — register routes during boot only.
Channel names are dot-separated tokens (e.g. campaigns.abc-123.donations.xyz-999). Subscribers can use two NATS-style wildcards to receive events across multiple concrete channels:
| Pattern | Matches |
|---|---|
campaigns.abc-123 |
exactly that channel |
campaigns.* |
campaigns.abc-123, campaigns.xyz-999 — one non-empty segment, nothing deeper |
campaigns.*.donations.* |
campaigns.abc.donations.xyz only — * always matches exactly one segment |
campaigns.> |
any channel starting with campaigns. (one or more segments) |
> |
everything published |
* may appear anywhere; > must be the trailing token. Published channel names must be concrete — wildcards in publish are rejected. Empty segments (campaigns..x) are rejected on both sides.
This lets pages scope their SSE stream to just what they need. Use an exact channel for a single-entity detail page, and a wildcard for a list or dashboard that should refresh on any change under a prefix:
class DonationPage < Sidereal::Page
# Only events for this specific donation
def channel_name = "campaigns.#{@campaign_id}.donations.#{@donation_id}"
end
class CampaignsListPage < Sidereal::Page
# Every campaign event and every donation event under every campaign
def channel_name = 'campaigns.>'
end
Pair this with a hierarchical channel_name block on the App to get routing “for free” from the channel name alone:
class DonationsApp < Sidereal::App
channel_name do |msg|
if msg.type.start_with?('donations.')
"campaigns.#{msg.payload.campaign_id}.donations.#{msg.payload.donation_id}"
else
"campaigns.#{msg.payload.campaign_id}"
end
end
end
Define inline components as separated classes (or nested classes) for partial re-renders:
class TodoPage < Sidereal::Page
class TodoList < Sidereal::Components::BaseComponent
def initialize(todos)
@todos = todos
end
def view_template
div(id: 'todos') do
@todos.each do |todo|
li { todo.title }
end
end
end
end
end
The command helper renders a form wired to POST /commands via Datastar. It handles hidden fields, AJAX submission, and server-side validation error display automatically.
def view_template
command AddTodo, class: 'add-form' do |f|
f.text_field :title, placeholder: 'What needs to be done?'
button(type: :submit) { 'Add' }
end
# Hidden payload fields (not shown to the user)
command RemoveTodo do |f|
f.payload_fields(todo_id: todo.todo_id)
button(type: :submit) { 'Remove' }
end
end
Field helpers: text_field, text_area, number_field, date_field, check_box, and payload_fields for values the user doesn’t edit. Pass a message instance instead of a class to prefill the form. Values are converted to and from the payload’s declared types on the way in and out — see Serialization.
Define a layout by subclassing Sidereal::Components::Layout. The base class overrides head and body to automatically inject the necessary Datastar wiring:
head — appends the Datastar JS script tag after your content.body — adds page signals (page_key, params) to the data attribute and appends the SSE init div at the end.class AppLayout < Sidereal::Components::Layout
def view_template
doctype
html do
head do
meta(name: 'viewport', content: 'width=device-width, initial-scale=1.0')
title { 'My App' }
end
body do
div(class: 'page') do
render page # renders the current page component
end
end
end
end
end
You can pass additional data attributes and signals to body. Extra signals are merged with the default page signals:
body(data: { class: 'app', signals: { theme: 'dark' } }) do
render page
end
Set the layout in your App:
class MyApp < Sidereal::App
layout AppLayout
# ...
end
A BasicLayout with reset CSS and form styling is provided by default if no layout is specified.
Chain .at(time) or .in(duration) (aliases) on a dispatch call to defer processing of a command (or event) until a future time. Three accepted forms:
| Form | Example | Resolution |
|---|---|---|
Time / DateTime |
.at(Time.now + 86400) |
Absolute target. |
Integer |
.in(3600) |
Seconds added to Time.now. |
Duration String |
.at('5m'), .in('PT1H30M') |
Parsed via Fugit.parse_duration, added to Time.now. |
command PlaceOrder do |cmd|
ORDERS[cmd.id] = cmd.payload.to_h
# Run an hour from now — Integer form
dispatch(SendReminder, order_id: cmd.id).in(3600)
# Run 30 minutes from now — Fugit duration String
dispatch(NudgeUser, order_id: cmd.id).in('30m')
# Run at a specific instant — Time form
dispatch(ExpireOrder, order_id: cmd.id).at(Time.now + 86400)
# ISO8601 durations also work
dispatch(SendDigest, order_id: cmd.id).in('PT1H')
end
Resolved targets earlier than the message’s created_at raise Sourced::Message::PastMessageDateError — including negative integers (.in(-60)) and durations that resolve to the past.
The dispatched message is appended to the store with its created_at set to the resolved target. Stores that support scheduled delivery hold the message back and only deliver it once that time has passed:
| Store | Scheduled delivery |
|---|---|
Store::FileSystem |
Honored — future-dated messages are written to a scheduled/ directory and promoted by a background fiber when due. |
Store::Memory |
Ignored — messages are delivered immediately regardless of created_at. Use FileSystem when you need scheduling. |
Scheduling does not propagate across correlation: an event dispatched downstream of a scheduled command runs at its own created_at (i.e. immediately), not at the source’s future time.
App.schedule registers a sequence of moments in time where commands should fire. The Scheduler is a leader-only fiber that, on each tick, appends commands to the same store the rest of your app uses — so schedule handlers run on the worker pool, in parallel with everything else, with the same retry / dead-letter machinery.
class MyApp < Sidereal::App
schedule 'Daily cleanup', '5 0 * * *' do |cmd|
# Runs every day at 00:05.
# `cmd` is the auto-generated command, materialised as
# MyApp::Commander::Schedules::SchedDailyCleanup0Step0.
dispatch SweepStaleOrders, older_than: '1d'
end
end
The schedule name ('Daily cleanup') is mandatory — it shows up in dead-letter sidecars, dashboards, and cmd.metadata[:schedule_name] so handlers and reactions can identify the source.
A step’s expression is anything Fugit.parse accepts, plus stdlib Time / DateTime instances:
| Kind | Example | Behaviour |
|---|---|---|
| Specific datetime | '2026-12-31T10:00:00' |
Fires once at that instant. |
Time instance |
Time.now + 60 |
Fires once at that instant (coerced internally). |
| Cron (5- or 6-field) | '5 0 * * *', '*/5 * * * *' |
Recurring at every cron match. |
| Natural language | 'every 3 seconds' |
Recurring. |
| Duration | '10d', '1h30m', 'P12Y12M' |
Fires once at previous concrete time + duration. |
at callsFor workflows that don’t fit a single step, drop the second positional and use the inner DSL — each at call appends a step to the schedule:
schedule 'Flash sale campaign' do
at '2026-05-10T10:00:00' do |cmd|
# Fires once at this exact moment.
dispatch OpenSale, sale_id: 'flash-2026'
end
at 'every day at 9am' do |cmd|
# Recurring — fires daily until the next concrete step.
dispatch SendDailyReminders
end
at '10d' do |cmd|
# Fires once at "previous concrete + 10 days".
# This concrete time also closes the recurring step above.
dispatch CloseSale, sale_id: 'flash-2026'
end
end
Steps are validated at registration:
'10d' is “10 days after the '2026-05-10T10:00:00' opening step”, not “10 days after the recurring started”. The resolved time also closes the recurring window.Drop the block / class for a specific or duration step to declare a bound-only marker — anchors the timeline without dispatching anything. Useful as a starting boundary for a following recurring step, or as a closing boundary for a preceding one:
schedule 'Office hours' do
at '2026-05-10T09:00:00' # marker — opens the window, no command
at 'every 5 minutes' do |cmd|
dispatch HealthCheck
end
at '2026-05-10T17:00:00' # marker — closes the recurring, no command
end
Block-less markers only work for specific or duration steps. A recurring step without a block would fire nothing on every match — meaningless — so it raises at registration.
By default, each at block generates a per-step command class under <HostCommander>::Schedules (e.g. MyApp::Commander::Schedules::SchedDailyCleanup0Step0 — Sched<CamelName><ScheduleId>Step<StepIndex>). For steps that should dispatch a domain command you’ve already defined elsewhere, pass the class + payload kwargs instead of a block:
# Define explicit commands and handlers
StartCampaign = Sidereal::Message.define('myapp.start_campaign')
command StartCampaign do |cmd|
# do something here
end
# Now just define time-based triggers for your own commands
schedule 'Flash sale campaign' do
at '2026-05-10T10:00:00', StartCampaign
at 'every day at 9am', SendEmails, sender: 'acme@company.org'
at '10d', EndCampaign
end
In the explicit form the macro generates no class and defines no handler — it just passes the class and payload through to the Scheduler, which dispatches SendEmails.parse(payload: { sender: 'acme@company.org' }, metadata: { ... }) on every fire. You’re responsible for having a command SendEmails do |cmd| ... end registered.
You can mix block and explicit forms freely across steps in the same schedule.
The Scheduler stamps these metadata keys on every dispatched command:
{
producer: "Schedule #0 'Flash sale campaign' step #1 (every day at 9am)",
schedule_name: "Flash sale campaign"
}
The producer label includes the schedule’s registration index, name, step index, and the step’s own expression — so dead-letter sidecars and dashboards can pinpoint which step fired. Both keys propagate to downstream commands via Message#correlate, so anything dispatched from inside a schedule handler carries them automatically.
The Scheduler ticks only on the process that holds Sidereal.elector. With the default Elector::AlwaysLeader (single-process apps) every process is leader; with Elector::FileSystem only one process per host runs the tick fiber. The dispatched commands then flow through the normal store, so any worker fiber on any process can pick them up — schedule handlers are not pinned to the leader.
A few caveats worth knowing:
crond’s no-catch-up behaviour.(@last_tick_at, now] — once the boundary moves into the past it can’t be in any future window.at '5m', X), it resolves to boot + 5m. Across a leader handoff each leader has its own boot time, so duration-anchored first steps drift; anchor with a specific datetime if stability matters.Sidereal::App is a subclass ofSidereal::Router , which is a standalone Rack-compatible router with a Sinatra-style DSL and trie-based dispatch. It can be used independently of the full Sidereal app framework.
Route blocks are evaluated in the context of a router instance, with access to request, response, and helper methods like body, status, headers, halt, and redirect.
class MyRouter < Sidereal::Router
get '/' do
body 'hello'
end
get '/items/:id' do |id:|
body "item #{id}"
end
post '/items' do
status 201
body 'created'
end
redirect '/old-path', '/new-path'
end
# config.ru
run MyRouter
Any object responding to #call(request, response, params) can be used as a handler. Callable handlers can either modify the response object or return a raw Rack triplet ([status, headers, body]).
class ShowItem
def call(request, response, params)
response.body = ["item #{params[:id]}"]
end
end
class MyRouter < Sidereal::Router
get '/items/:id', ShowItem.new
# Lambdas work too
get '/health', ->(req, resp, params) { [200, {}, ['ok']] }
end
Run logic before every matched route. Use halt to short-circuit.
class MyRouter < Sidereal::Router
before do
halt 401, 'unauthorized' unless session[:user_id]
end
get '/dashboard' do
body 'welcome'
end
end
class MyRouter < Sidereal::Router
session secret: ENV.fetch('SESSION_SECRET')
post '/login' do
session[:user_id] = request.params['user_id']
body 'logged in'
end
get '/profile' do
body "user: #{session[:user_id]}"
end
end
halt immediately stops request processing and returns a response.
halt 422 # status only
halt 200, 'hello' # status + body
halt 301, 'Location' => '/new-path' # status + headers
halt 200, { 'X-Custom' => '1' }, 'ok' # status + headers + body
redirect is a shorthand for halting with a Location header:
redirect '/new-path' # 301 by default
redirect '/new-path', status: 302 # temporary redirect
Sidereal.dependencies is a container for the resources an app sets up at boot — a database connection, an API client — and for Sidereal’s own services. Register them in config/dependencies/*.rb, loaded from boot.rb:
# config/dependencies/db.rb
require 'sequel'
Sidereal.dependencies.register!('db') do
Sequel.sqlite(ENV.fetch('DB_PATH'))
end.teardown do |db|
db.disconnect
end
# config/dependencies/repos.rb
Sidereal.dependencies.register!('repos.orders', ['db', 'logger']) do |db, logger|
OrdersRepo.new(db, logger:)
end
# boot.rb
Dir[File.join(__dir__, 'config/dependencies/*.rb')].sort.each { |f| require f }
register! registers a singleton, built once per process and reused. register registers a transient dependency, built again on every lookup..teardown { |value| ... } releases a singleton at shutdown.requires belong at the top of the file, outside the block: they run as the file loads, which is shared by every worker when the app is preloaded.Look a value up with Sidereal.dependencies['repos.orders'], which builds 'db' and 'logger' first if they haven’t been built.
Sidereal::Host#start calls Sidereal.dependencies.build! before anything else, in every process, leader or follower. It checks the graph — a missing key, a cycle, or a class injecting an unregistered key fails the boot — builds every singleton in dependency order, and locks the container against further registration. Anything a process needs set up, whatever its role, belongs here. Building happens while the channel and exception registries are still open, so a dependency’s block can register an exception subscriber:
Sidereal.dependencies.register!('error_reporter') do
client = ErrorReporter::Client.new(ENV.fetch('ERROR_REPORTER_KEY'))
Sidereal.exceptions.on_failure { |report| client.notify(report.exception) }
client
end
Sidereal::Host#stop calls teardown, which runs the teardowns of the singletons built in that process, dependents before their dependencies.
Outside a host — specs, rake tasks, a console — lookups build what they need on demand, or call Sidereal.dependencies.build! to build everything.
Values are never shared across processes. The Falcon controller, which forks the workers, refuses to build anything (Sidereal::Dependencies::ForkError), and a value built in one process raises the same error if another process tries to use it.
args adds keyword arguments to a class’s constructor, defaulting to the container’s values, with a reader for each:
class OrdersProjector < Sourced::Projector::StateStored
include Sidereal.dependencies.args('db')
sync do |state:, **|
db[:orders].insert_conflict(:replace).insert(state)
end
end
OrdersProjector.new(partition_values) # db: Sidereal.dependencies['db']
OrdersProjector.new(partition_values, db: test_db)
'sourced.store' → store:). A hash renames it: args('sourced.store' => 'events').includes add up, and every other argument reaches the class’s own initialize untouched — it doesn’t need to call super — so classes that frameworks instantiate themselves, like Sourced’s reactors, work unchanged.Sidereal::Deps is class-level shorthand for the same thing — dep is include Sidereal.dependencies.args(...):
class OrdersProjector < Sourced::Projector::StateStored
extend Sidereal::Deps
dep :db
dep 'sourced.store' => 'events'
end
A name the class already has is refused, since the reader would replace it — alias it instead (dep 'store' => 'orders_store').
Commanders extend it already, and with the Sourced integration loaded so do Sourced::Decider and Sourced::Projector, so their handlers can use what they declare:
class Orders < Sidereal::Commander
dep 'repos.orders' => 'orders'
command PlaceOrder do |cmd|
orders.insert(cmd.payload)
end
end
On an App, dep reaches both kinds of handler: handle blocks, which run on the app for each request, and command blocks, which run on its commander:
class ShopApp < Sidereal::App
dep 'repos.orders' => 'orders'
handle PlaceOrder do |cmd|
halt 422 if orders.duplicate?(cmd.payload)
dispatch cmd
end
command PlaceOrder do |cmd|
orders.insert(cmd.payload)
end
end
Commanders added with commands declare their own, and an App subclass has a commander of its own, so it declares again what its command blocks use.
The store, pubsub and elector are registered as 'sidereal.store', 'sidereal.pubsub' and 'sidereal.elector'. Sidereal.store and friends read them, and c.store = ... and the other setters replace them. To replace one with something built from other dependencies, register it with override: true:
Sidereal.dependencies.register!('sidereal.store', ['db'], override: true) do |db|
MyStore.new(db)
end
A key registers once: registering it again without override: true raises. An override is refused once the key, or anything built from it, has been resolved (whatever holds the old value would keep it), and every registration is refused after build!.
Sidereal is designed to run on Falcon, which provides the async fiber runtime needed for SSE streaming and concurrent command processing.
Create a falcon.rb file. It requires only the environment: each worker loads the app itself through config.ru (see Preload vs lazy loading).
#!/usr/bin/env falcon-host
# frozen_string_literal: true
require 'sidereal/falcon/environment'
service "my-app" do
include Sidereal::Falcon::Environment
include Falcon::Environment::Rackup
url "http://localhost:9292"
count 1
end
Run with:
bundle exec falcon host
Sidereal.configure do |c|
c.workers = 3 # number of worker fibers processing commands
end
Falcon forks worker processes (count, which defaults to the CPU count). Where the app loads decides what the workers share, and what a zero-downtime restart picks up. Whichever you choose, every worker builds its own connections: Sidereal builds dependencies in each worker, when it boots.
Lazy (default). falcon.rb requires only the environment, and each worker loads config.ru → boot.rb in its own process. Nothing of the app is shared across workers, and Falcon’s zero-downtime restart (SIGHUP), which forks a fresh set of workers, loads the new code from disk — so HUP deploys work.
# falcon.rb — lazy: the app loads per worker via config.ru
require 'sidereal/falcon/environment'
service "my-app" do
include Sidereal::Falcon::Environment
include Falcon::Environment::Rackup
url "http://localhost:9292"
end
# config.ru — loads boot.rb inside each worker
require_relative 'boot'
run MyApp
Preload gems. Gems in a :preload Bundler group are required by Falcon’s controller, which then compacts its heap (Process.warmup), so workers forked from it share them copy-on-write. App code still loads per worker, and HUP deploys still work. As async-service implements it, the warm-up runs after the first workers have started, so it is the workers forked by later restarts that benefit.
Preload the app (opt-in). preload "boot.rb" in the service block loads the app once in the controller, before forking: code, classes and compiled codecs are shared copy-on-write, and workers boot faster. Two things to know:
HUP forks new workers from the controller, which still holds the code it loaded at start. Restart falcon host (or switch instances in front of it) to pick up new code.Sidereal::Dependencies::ForkError instead of handing every worker the same connection. A connection opened outside the container can’t be caught, though: keep them in dependencies rather than, say, Sourced.configure { |c| c.store = Sequel.sqlite(...) } or a Sequel::Model that reads its schema when defined. Preload through preload, not by requiring the app from falcon.rb, which runs before that check is in place.service "my-app" do
include Sidereal::Falcon::Environment
include Falcon::Environment::Rackup
url "http://localhost:9292"
preload "boot.rb" # everything the app needs; config.ru then just requires it
end
For an in-memory backend the distinction hardly matters; it bites when a backend holds a real, fork-unsafe connection.
The store, pubsub, elector and dispatcher are configurable. By default Sidereal uses in-memory implementations, but you can swap them out:
Sidereal.configure do |c|
c.store = MyCustomStore.new # default: Sidereal::Store::Memory
c.pubsub = MyCustomPubSub.new # default: Sidereal::PubSub::Memory
c.dispatcher = MyDispatcherClass # default: Sidereal::Dispatcher (a class, not an instance)
end
A custom store must respond to #append(message). A custom dispatcher must respond to .start(task) (class-level) and #stop. The store, pubsub and elector are dependencies ('sidereal.store' and so on) — the setters above replace them, and a replacement that needs other dependencies, such as the app’s database, is registered directly.
Which process runs the dispatcher. c.dispatcher_process is :all by default: every process starts a dispatcher at boot. Set it to :leader and only the process holding Sidereal.elector starts one — the same rule the Scheduler follows. The dispatcher is stopped if that process is demoted, and the next leader starts its own. Web requests keep appending commands from every process; only the consuming side is pinned. That is the right shape for a backend that serializes writers (Sourced on SQLite, where the Sourced integration sets it for you): reads scale across workers while handler and projection writes come from one.
Sidereal.configure do |c|
c.use_file_system! # a cross-process elector is what makes :leader meaningful
c.dispatcher_process = :leader
end
With the default Elector::AlwaysLeader every process is leader, so the two modes coincide. A dispatcher that fails to start on a later promotion (after a failover) is logged by the elector’s callback guard rather than failing the boot, since promotion happens after boot.
Multi-process shortcut. c.use_file_system! switches the store, pubsub, and elector to their filesystem / unix-socket implementations in one call — the combination needed to run across multiple Falcon workers on one host (a shared on-disk queue, a unix-socket pubsub broker, and file-lock leader election). Files and the socket live under dir: (default ./storage, relative to the working directory). Override any individual collaborator afterward:
Sidereal.configure do |c|
c.use_file_system! # FS store + unix-socket pubsub + file-lock elector
c.store = Sourced.config.store # ...e.g. keep Sourced's store, but the rest stays
end
Integrations. Backends that provide several collaborators at once (e.g. a store + dispatcher pair, plus bridging) ship as integrations, applied with c.use(SomeIntegration, **opts) — use_file_system! is itself one. See Using Sourced as a backend for the canonical example.
Sidereal::Store::FileSystem is a built-in durable store that survives process restarts and lets multiple worker processes on the same host share a queue. It also honors scheduled commands, unlike the default in-memory store.
It isn’t autoloaded — require it explicitly, then point Sidereal.configure at an instance:
require 'sidereal'
require 'sidereal/store/file_system'
Sidereal.configure do |c|
c.store = Sidereal::Store::FileSystem.new(root: 'storage/store')
end
File bodies are one JSON document per message, written by the shared transport codec — see Serialization for the wire shape and for adding encoders for your own payload types.
The store creates five sibling directories under root/: tmp/, ready/, scheduled/, processing/, and dead/. Producers append by atomic-renaming from tmp/ into ready/ (or scheduled/ for future-dated messages). A poller fiber claims into processing/; a scheduler fiber promotes due files from scheduled/ to ready/; a sweeper recovers anything left in processing/ by a crashed worker. Permanently-failed messages land in dead/ along with a <f>.error.json sidecar — see Failure handling.
Constructor options:
| Option | Default | Description |
|---|---|---|
root: |
'tmp/sidereal-store' |
Directory holding the four subdirs. Must be on a single local filesystem (atomic rename is unreliable across NFS). |
poll_interval: |
0.1 |
Seconds the poller sleeps when ready/ is empty. |
scheduler_interval: |
1.0 |
Seconds between scans of scheduled/ for due messages. Sub-second granularity is not provided. |
sweep_interval: |
60 |
Seconds between sweeps of stale processing/ files. |
stale_threshold: |
300 |
A processing/ file older than this (or owned by a dead PID) is treated as abandoned and renamed back to ready/. |
max_in_flight: |
50 |
Bound on the in-process queue between the poller and worker fibers. When handlers fall behind, the queue blocks and disk becomes the buffer. |
At-least-once delivery: a crash mid-handling causes the message to be re-claimed once the sweeper recovers the abandoned processing/ file. Handlers must be idempotent.
Single-machine only: the design relies on POSIX atomic rename, which is unreliable across networked filesystems like NFS. Use a different store if you need to fan workers out across hosts.
When a command handler raises, the dispatcher calls Commander.on_error(exception, message, meta) and uses the returned value to decide what to do next:
| Return value | Effect |
|---|---|
Sidereal::Store::Result::Retry.new(at: time) |
re-schedule for another attempt at time |
Sidereal::Store::Result::Fail.new(error: exception) |
give up — dead-letter the message |
Sidereal::Store::Result::Ack |
swallow silently — drop the message |
The default policy retries with exponential backoff (2 ** meta.retry_count seconds) up to Sidereal::Commander::DEFAULT_MAX_ATTEMPTS attempts, then fails. Override per-commander:
class MyApp < Sidereal::App
commands do
def self.on_error(exception, message, meta)
case exception
when MyDomain::Invalid
Sidereal::Store::Result::Fail.new(error: exception) # bail immediately
when Net::Timeout
Sidereal::Store::Result::Retry.new(at: Time.now + (5 * meta.retry_count))
else
super # fall back to the default policy
end
end
end
end
meta.retry_count starts at 1 and increments on each retry. meta.first_appended_at is preserved across retries — useful for “give up after N hours regardless of attempt count” policies.
Sidereal::Store::FileSystem — Retry renames the message into scheduled/ with a bumped retry_count and the new not_before_ns; the body stays untouched (commanders cannot mutate the message between attempts). Fail writes a sidecar dead/<f>.error.json with the exception class/message/backtrace, then renames the message into dead/. The sweeper does not touch dead/ — those messages are terminal until you act on them manually.Sidereal::Store::Memory — Retry and Fail log at WARN level and ack/drop the message. The in-memory store has no scheduling or dead-letter primitives.Requeueing dead messages. Once you’ve fixed the underlying cause of failure, Sidereal::Store::FileSystem#requeue(filename) moves a dead-lettered message back into ready/. The new filename has retry_count reset to 1 and not_before_ns set to now (immediately ready); first_append_ns is preserved so age-based diagnostics retain the lineage. The .error.json sidecar is removed.
store = Sidereal::Store::FileSystem.new(root: 'storage/store')
store.requeue('1762000000-1761000000-3-12345-abcdef.json')
# => "<root>/ready/<new-filename>.json"
Path components in the input are stripped via File.basename, so 'abc.json', 'dead/abc.json', and '/abs/dead/abc.json' are all equivalent — the file is always resolved against the store’s configured dead/ directory. Missing files raise ArgumentError.
At-least-once delivery still applies: a worker crash before Retry/Fail is acted on leaves the file in processing/ for the sweeper to recover, which re-runs the handler. Handlers must be idempotent.
Every retry or terminal failure decision a backend makes is reported through Sidereal.exceptions, a process-global subscriber registry. Subscribers receive a small ExceptionReport value:
ExceptionReport = Data.define(:kind, :exception, :message, :retry_count, :retry_at)
# kind: :retry | :failure
# exception: the raw StandardError instance
# message: the failed Sidereal::Message (typically a command)
# retry_count: 1-indexed attempt number that just failed
# retry_at: Time of the next attempt (nil on :failure)
Register subscribers during boot — APM hooks, structured loggers, anything you want notified:
Sidereal.exceptions.on_failure do |report|
Sentry.capture_exception(report.exception, extra: report.message.payload.to_h)
end
Sidereal.exceptions.on_retry do |report|
StatsD.increment('handler.retry', tags: ["command:#{report.message.class.type}"])
end
Backends call into the registry from inside their retry/fail policy:
Sidereal.exceptions.report_retry(exception:, message:, retry_count:, retry_at:)
Sidereal.exceptions.report_failure(exception:, message:, retry_count:)
Sidereal::Dispatcher does this automatically — its dispatch_notification is the only place that calls these methods today. Other dispatchers (Sourced, custom) wire their own retry/fail callbacks the same way.
Sidereal.exceptions ships with a default subscriber pair pre-installed via Sidereal::Exceptions.with_default_publisher. Each one builds the corresponding Sidereal::System::Notify* from the report and broadcasts it on the failed message’s channel:
Sidereal::System::NotifyRetry — payload: command_type, command_id, command_payload, retry_count, retry_at (ISO8601), error_class, error_message, backtrace.Sidereal::System::NotifyFailure — same payload minus retry_at.Both inherit from Sidereal::System::Notification (a marker base). The default publisher publishes them via Sidereal.pubsub.publish(Sidereal.channels.for(failed_command), notify), where Sidereal.channels ships with pre-installed source-channel bypass routes so the resolution lands on the originating command’s channel without any user-supplied resolver having to know about system messages.
Pages render the toasts via the default reactions in Sidereal::Page:
on Sidereal::System::NotifyFailure do |evt|
browser.patch_elements Sidereal::Components::SystemNotifyFailure.new(evt),
mode: 'prepend', selector: '#sidereal-sysnotify-stack'
end
Override on your own page subclass to render a custom UI:
class TodoPage < Sidereal::Page
on Sidereal::System::NotifyFailure do |evt|
browser.patch_elements MyErrorBanner.new(evt)
end
end
The dispatcher’s report-call site short-circuits when the failing message is itself a Sidereal::System::Notification. So a buggy on_failure subscriber whose own exception cycles back into the worker doesn’t trigger a fresh report-and-fan-out. Subscriber exceptions are also caught by the registry and logged via Console.error — a single broken subscriber never tears down the worker fiber or prevents later subscribers from firing.
The base Sidereal::Page ships with default reactions that render Sidereal::Components::SystemNotifyRetry (amber) or SystemNotifyFailure (red) toasts and prepend them into a fixed-position stack at the top-right of the page (#sidereal-sysnotify-stack, supplied by the base layout’s sidereal_foot). Each toast shows the command type, error class/message, attempt count, retry time, and a collapsible backtrace; they slide in/out, are dismissable, and carry their own inline <style> so they don’t depend on host CSS.
The default reactions fire in any environment for now. Override on(NotifyRetry) / on(NotifyFailure) on your page to render a custom UI, or use Sidereal::Exceptions.new (without the .with_default_publisher factory) and inject it via the dispatcher’s exceptions: kwarg to suppress the publish entirely for headless deployments.
module Sidereal
module System
NotifyDeprecated = Notification.define('sidereal.system.notify_deprecated') do
attribute :command_type, Sidereal::Types::String
attribute :reason, Sidereal::Types::String
end
end
end
Defining via Notification.define(...) registers it under the Notification registry; the dispatcher’s loop prevention keys off is_a?(Notification) and picks up the new class automatically. You’ll still need to:
:source_channel-bypass route for it on Sidereal.channels (mirroring the bypass installed for NotifyRetry/NotifyFailure) so it reaches the originating page’s SSE channel;Sidereal::Exceptions.build_notification (or register a custom subscriber that handles the new kind) so reports are translated into the new message;Page.on(...) reaction;A command crosses two boundaries with very different shapes. It goes over the wire to a store or a pub/sub socket, where it must become bytes; and it goes through an HTML form, where every value — a date, a number, a checkbox — is a String in both directions.
Sidereal compiles a codec for each, from your command payload schemas. You never call either one directly. The point is that you declare payload attributes in the types you actually want to work with, and each boundary translates:
| Codec | Crosses | Encodes | Compiled by |
|---|---|---|---|
Sourced::Message::JSONCodec |
Store::FileSystem file bodies, PubSub::Unix frames |
the whole message, envelope included | Store#start / #append, PubSub#start / #publish |
Sidereal::FormsCodec |
POST /commands params, <input value="..."> |
the payload alone, every scalar a String | App.handle |
Both are built on Plumb’s codecs — Plumb::Codec::JSON and Plumb::Codec::Forms — which rewrite a schema into a decoder/encoder pair by resolving an encoder for every leaf type.
A browser submits seats=30 as the String "30", published as "1", and a date as "2026-09-01". Without a codec you would either coerce by hand in every handler, or write your schemas in lax types and lose the guarantee. Instead, declare what you mean:
BookCourse = Sidereal::Message.define('courses.book') do
attribute :course_name, Sidereal::Types::String.present
attribute :seats, Integer
attribute :starts_on, Date
attribute :published, Sidereal::Types::Boolean
end
class CoursesApp < Sidereal::App
handle BookCourse
command BookCourse do |cmd|
cmd.payload.seats # => 30 (Integer)
cmd.payload.starts_on # => #<Date 2026-09-01>
cmd.payload.published # => true (TrueClass)
# so this just works, with nothing parsed by hand
dispatch Reminder.at(cmd.payload.starts_on - 7) if cmd.payload.published
end
end
Only commands that are web-facing (via .handle) are form-decoded. A command registered only with command is never reachable from a form and is never compiled.
What Plumb::Codec::Forms knows out of the box:
| Attribute type | Accepts from a form | Renders back as |
|---|---|---|
Types::String |
any string | itself |
Types::Integer |
"30", "-4" |
"30" |
Types::Float / Types::Decimal |
"1.5", "9.99", "1e3" |
"1.5" |
Types::Boolean |
"true"/"1", "false"/"0" (case-insensitive) |
"true" / "false" |
Types::Date |
"2026-09-01" |
"2026-09-01" |
Types::Time |
ISO 8601 | "2026-09-01T10:00:00.000000+01:00" |
Types::Symbol |
any string | itself |
Types::URI::Generic / ::HTTP / ::File |
an RFC 3986 URI | itself |
anything .nullable |
"", or an absent field |
"" |
The same translation runs backwards, so command accepts a message instance as well as a class. A class renders a blank form; an instance prefills each field:
# A class — every field renders empty
command BookCourse do |f|
f.text_field :course_name # <input type="text" name="command[payload][course_name]">
f.text_area :description # <textarea name="command[payload][description]"></textarea>
f.number_field :seats
f.date_field :starts_on
f.check_box :published
end
# An instance — every set attribute renders its encoded value
command BookCourse.new(payload: {
course_name: 'Ruby 101', seats: 30,
starts_on: Date.new(2026, 9, 1), published: true
}) do |f|
f.text_field :course_name # <input type="text" ... value="Ruby 101">
f.date_field :starts_on # <input type="date" ... value="2026-09-01">
f.check_box :published # checked
end
An attribute that is unset renders no value attribute at all, so the same form definition serves both cases — including the one in between, a command half-filled from a previous attempt, where the attributes that are set render and the rest come out blank.
The payload is encoded once per render, not once per field, and per key rather than all-or-nothing. That is what lets a blank or partial command render at all: a strict conversion would reject one outright.
Encoding also happens only at an input’s value=. The command object itself keeps its Ruby values, so logic inside the form block sees what you’d expect:
command course_cmd do |f|
f.date_field :starts_on
# a real Date and a real boolean — not "2026-09-01" and "1"
p { "Starts in #{(f.command.payload.starts_on - Date.today).to_i} days" }
f.check_box :published unless f.command.payload.published
end
check_box renders a hidden 0 alongside the checkbox, because an unchecked box submits nothing at all. Rack keeps the last value for a repeated name, so a checked box sends "1" and an unchecked one "0".
payload_fields carries values the command doesn’t hold — an id from a loop variable, a preset amount. Those are encoded as attributes of that command, so a hidden field and a visible one for the same attribute always agree, and a key the payload doesn’t declare raises rather than rendering an empty input.
When decoding fails, the result is a flat {attribute => message} hash, which POST /commands streams straight back to the offending field over SSE — see Command forms:
command[payload][seats]=lots
# => the "seats" field gets: Must match /\A-?\d+\z/
Two behaviours worth knowing:
Types::Integer.default(0) plus a blank input errors. .default fires for an absent key, and "" is a present value that no Integer encoder accepts. Use Types::Integer.nullable for optional numeric fields.Types::Boolean reads Must match /\Atrue\z/i, Must be equal to 1, ....Stores and pub/sub serialize the whole message — envelope included — because a file body or a socket frame has nowhere else to put an id, a created_at or a correlation chain. That is Sidereal.message_codec, shared by Store::FileSystem and PubSub::Unix:
{
"id": "97b72e83-c27c-4a0f-b8e7-19abbea9f70e",
"causation_id": "97b72e83-c27c-4a0f-b8e7-19abbea9f70e",
"correlation_id": "97b72e83-c27c-4a0f-b8e7-19abbea9f70e",
"created_at": "2026-08-10T19:51:35.032711+01:00",
"metadata": {},
"type": "courses.book",
"payload": {
"course_name": "Ruby 101",
"seats": 30,
"starts_on": "2026-09-01",
"published": true
}
}
Note that the two formats disagree, correctly, about the same schema: seats is a JSON number here and the String "30" in a form, and starts_on is an ISO date string in both but a Date at rest in Ruby. Each codec keeps its own compiled pair per message class, which is what makes that possible.
Nothing compiles on first use. Each transport calls compile! when it starts (and again on the first write, since Sidereal.dispatch! from a CLI can append with no dispatcher running), so a schema the format cannot represent fails at boot rather than on the message that happens to carry it.
The web boundary never reads the envelope. A form supplies command[type] and the payload; id, created_at, metadata and the correlation chain are built server-side, so a request cannot date a command into the future and have the store schedule it.
Sooner or later a payload carries something neither format knows — a Money, a Coordinate, a domain enum. Declare an Encoder for it and register it on each codec it will cross. Both are separate registries: teaching one does not teach the other.
An encoder is a class declaring Input => Output plus the two conversions. Output is your Ruby type; Input is the shape the format can carry:
# money.rb
Money = Data.define(:cents, :currency) do
def self.euros(units) = new(cents: units * 100, currency: 'EUR')
def to_s = "€#{cents / 100}"
end
# A form field is a String and nothing else, so both parts are packed into one.
# The `Types::` form is needed here because it is a *refinement* — String, but
# only strings matching that pattern.
class MoneyFormsEncoder < Plumb::Encoder[
Plumb::Types::String[/\A\d+ [A-Z]{3}\z/] => Money
]
def encode(money) = "#{money.cents} #{money.currency}"
def decode(str)
cents, currency = str.split
Money.new(cents: cents.to_i, currency:)
end
end
# JSON has objects, so the parts can stay addressable on the wire. A plain class
# is enough where no refinement is involved.
class MoneyJSONEncoder < Plumb::Encoder[
Plumb::Types::Hash[cents: Integer, currency: String] => Money
]
def encode(money) = { cents: money.cents, currency: money.currency }
def decode(hash) = Money.new(cents: hash[:cents], currency: hash[:currency])
end
Plumb::Codec::Forms.encoder(MoneyFormsEncoder)
Plumb::Codec::JSON.encoder(MoneyJSONEncoder)
The two need not agree on a shape, and here they deliberately don’t. One Ruby type reaches each wire in the form that wire can carry:
SelectAmount = Sidereal::Message.define('donations.select_amount') do
attribute :amount, Sidereal::Types::Any[Money]
end
hidden form field value="3000 EUR"
store file / frame "amount": { "cents": 3000, "currency": "EUR" }
command handler Money[cents: 3000, currency: "EUR"]
Rendering and submitting both go through the encoder, so a preset-amount button is just:
command SelectAmount, key: amount.cents do |f|
f.payload_fields(amount:) # => <input type="hidden" value="3000 EUR">
button(type: :submit) { amount.to_s }
end
Register encoders at load time, before any message type is defined. Every compile walks the whole message registry, so a format missing an encoder for a type any message uses cannot compile at all. require the file at the top of your boot sequence — see examples/donations1/money.rb.
If you get it wrong you find out immediately, and the error names the attribute path:
cannot apply Plumb::Codec::Forms[...] (decode) to Booking::Payload:
field `window` (Range[Integer]) matches no encoder and is not covered by its
noop types. Register an encoder for it, or declare it with .noop.
A type can be representable in one format and not the other, and that is fine — it just means the command cannot be web-facing. Types::Range is the built-in example: Plumb::Codec::JSON encodes it as {from:, to:, exclusive:}, while Plumb::Codec::Forms deliberately does not register it, since a single form field has no sensible shape for it. Such a command serializes for transport, and raises at handle if you try to expose it to the browser.
sequenceDiagram
participant Browser
participant App
participant PubSub
participant Store
participant Worker
participant CommandHandler
Browser->>App: POST /commands
App->>App: Check handled_commands registry (404 if not exposed)
App->>App: Validate command
App->>Store: Append command (default handler)
App->>Browser: 200 OK
loop Worker fibers
Worker->>Store: Claim next command
Store->>Worker: Command
Worker->>CommandHandler: Handle command
CommandHandler->>Worker: Result(events, commands)
Worker->>PubSub: Publish events
Worker->>Store: Append new commands (if any)
end
Browser->>App: GET /updates (SSE)
App->>PubSub: Subscribe
PubSub->>App: Event
App->>App: Page reactions render HTML
App->>Browser: SSE patch (HTML fragments)
Browser->>Browser: Datastar morphs DOM
Commands are processed asynchronously by worker fibers. The browser never waits for command handling to complete – it submits the command and receives UI updates via the SSE stream as events are produced.
Sourced (“ccc” branch) is an event sourcing library with a persistent SQLite store, partition-aware consumer groups, and a signal-driven dispatcher. Sidereal can use it as a drop-in backend, replacing the in-memory store and dispatcher.
Require the Sourced integration, register the database as a dependency, then point Sidereal at Sourced’s store and dispatcher:
require 'sequel'
require 'sidereal'
require 'sidereal/integrations/sourced'
Sidereal.dependencies.register!('db') do
Sequel.sqlite('db/app.db')
end.teardown(&:disconnect)
Sidereal.configure do |c|
# Use file-system version of pubsub, elector
c.use_file_system!
# Use Sourced as message store and dispatcher, on the 'db' connection
c.use Sidereal::Integrations::Sourced, store: 'db'
end
c.use Sidereal::Integrations::Sourced wires several things for you:
Store + dispatcher — Sourced becomes Sidereal’s message store and dispatcher. Sidereal Commanders (command / handle) also register as Sourced reactors, so they run on the same runtime alongside your Deciders and Projectors. store: sets Sourced’s store in each worker: a dependency key (as above — the same connection your own classes can inject), or a callable:
c.use Sidereal::Integrations::Sourced, store: -> { Sequel.sqlite('db/app.db') }
Without store:, the integration uses whatever you configured with Sourced.configure — which opens its connection while the app loads, so avoid it if you preload the app.
Sidereal.exceptions, so the default error toasts appear and any on_retry / on_failure / on_fatal subscribers (e.g. an APM hook) fire. When Sourced is the dispatcher it owns retry/fail orchestration, so Sidereal’s automatic exception reporting doesn’t run — this bridge is what surfaces failures in the UI.c.dispatcher_process = :leader (see Custom backends), so only the elected process runs the Sourced runtime — commanders, deciders and projectors. SQLite serializes writers, so N runtimes on N workers would queue on each other; with one, the other workers serve pages and queries in parallel. Every worker still appends commands (a form post appends from whichever worker served it); it is the claiming, handling and projecting that runs in one place. Set c.dispatcher_process = :all after use to run a runtime on every worker again.Sidereal::Integrations::Sourced::Notifier, which carries those announcements over Sidereal.pubsub — with the unix-socket pubsub, an append on any worker wakes the leader immediately. Appends made outside an Async reactor (a rake task calling Sidereal.dispatch!) cannot reach the socket and fall back to the catch-up poll, which stays the safety net in every case.'sourced' dependency whose block sets Sourced’s store from store:, calls Sourced.setup! and recompiles the store’s message codec. Sidereal::Host#start builds it in every process before anything starts. Only the leader runs the runtime, but every worker appends, so every worker needs Sourced’s store ready: its connection open, its tables installed and its codec compiled against every message type the app defines (Sourced.configure compiles it earlier, while boot.rb loads, before the app’s types exist). Because the store is set up there and never while the app loads, each worker opens its own connection, preloaded or not. Sourced.setup! freezes Sourced’s configuration, so it runs once per process: the dependency is that one call. A process without a Host, such as a rake task, resolves it itself after loading the app: Sidereal.dependencies['sourced'].Multi-process:
use_file_system!(cross-process pubsub + file-lock election) is required whenever you run more than one worker — otherwise the in-process pubsub/elector can’t fan SSE updates across processes. If you start multiple workers with the default in-process subsystems, Sidereal refuses to boot with a loud error telling you to add it (see Running with Falcon).
With Sourced as the backend, use Sourced messages and Deciders instead of App.command:
# Define messages using Sourced's message class. AddTodo carries todo_id
# because the Decider partitions by it (generate it client-side, or stamp it
# in a `before_command` hook).
AddTodo = Sourced::Message.define('todos.add') do
attribute :todo_id, String
attribute :title, String
end
TodoAdded = Sourced::Message.define('todos.added') do
attribute :todo_id, String
attribute :title, String
end
# Define a Decider (replaces App.command for async processing)
class TodoDecider < Sourced::Decider
partition_by :todo_id
command AddTodo do |state, cmd|
event TodoAdded, todo_id: cmd.payload.todo_id, title: cmd.payload.title
end
end
Sourced.register(TodoDecider)
You don’t write any publishing code — that’s the auto-publish below.
The integration publishes reactor output to Sidereal’s PubSub for you, so Page reactions pick it up and stream DOM updates via SSE. Every path resolves its channel through your App’s channel_name resolver (Sidereal.channels.for(evt)), and a publish/resolver failure is reported to Sidereal.exceptions (it’s terminal — the store already committed).
after_sync into every Sourced::Decider subclass, so no per-reactor wiring is needed.Projectors publish a synthetic “projected” signal so Pages know to re-fetch their read model. You don’t define or publish it: the integration wraps the partition_by macro to auto-define a Projected event class (one attribute per partition key) and register an after_sync that publishes it after each committed batch:
class TodosProjector < Sourced::Projector::StateStored
partition_by :todo_id
evolve TodoDecider::TodoAdded do |state, evt|
# update the read model...
end
sync do |state:, **|
# persist the read model...
end
end
# => auto-defines TodosProjector::Projected (with a `todo_id` attribute),
# published after every batch via Sidereal.channels.for.
Sourced.register(TodosProjector)
Pages react on the generated signal:
on TodosProjector::Projected do |_evt|
browser.patch_elements load(params)
end
The signal’s payload is the projector’s full partition tuple, so multi-key partitions work too — partition_by(:student_id, :course_id) yields a two-attribute Projected, and your channel_name resolver routes it just like your domain events.
The signal is correlated, not synthetic: a batch may hold messages from several causal chains, so the integration publishes one Projected per distinct correlation_type in the batch, each correlated from the last message of that chain. A page that renders command CreateGame, or declares on CreateGame without a block, therefore reloads when the projector commits the batch holding GameCreated, with no reaction written for the signal. An explicit on MyProjector::Projected do |evt| ... end still fires for every batch regardless of root, which is what a lobby listing every game wants.
Sidereal Commanders are unaffected — they’re not Decider/Projector subclasses and keep publishing via their own path, so there’s no double-publish.
Pages and the App work the same way. handle exposes commands to the browser, and Page on reactions respond to the auto-published events:
class TodoPage < Sidereal::Page
path '/'
on TodoAdded do |evt|
browser.patch_elements load(params)
end
def self.load(_params, _ctx)
new(todos: TODOS.values)
end
# ...
end
class TodoApp < Sidereal::App
session secret: ENV.fetch('SESSION_SECRET')
channel_name { |msg| "todos.#{msg.payload.todo_id}" }
# Expose AddTodo to the browser — the default handler
# appends it to the Sourced store for async processing
handle AddTodo
page TodoPage
end
You can also process a command synchronously during the HTTP request using Sourced.handle!, which loads history, runs the Decider, appends events, and returns immediately. Unlike the async dispatcher path, handle! does not run the reactor’s after_sync, so the auto-publish doesn’t fire here — publish the returned events yourself if other connected browsers need them:
handle AddTodo do |cmd|
_cmd, _decider, events = Sourced.handle!(TodoDecider, cmd)
events.each { |evt| pubsub.publish(Sidereal.channels.for(evt), evt) }
browser.patch_elements TodoList.new(TODOS.values)
end
The standard Sidereal Falcon environment works with Sourced – it uses the configured dispatcher automatically:
#!/usr/bin/env falcon-host
require 'sidereal/falcon/environment'
service "my-app" do
include Sidereal::Falcon::Environment
include Falcon::Environment::Rackup
url "http://localhost:9292"
count 1
end
After checking out the repo, run bin/setup to install dependencies. Then, run bundle exec rspec to run the tests. You can also run bin/console for an interactive prompt.
Join us in the sidereal tag on the Ruby Users Forum.
Bug reports and pull requests are welcome on GitHub at https://github.com/ismasan/sidereal.
The gem is available as open source under the terms of the MIT License.