diff --git a/CHANGELOG.md b/CHANGELOG.md index 2df631a..c745ae6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,18 @@ # Changelog +## 0.6.0 - 2026-08-09 + +- Add `broadcast_payload`, an actor DSL for sending one personalized JSON state + payload over the actor stream a page already has open. The block runs once per + subscriber with that subscriber's authorization context, so private state + never crosses sessions. Payloads carry actor identity and the monotonic state + revision, and both the channel and the browser drop stale revisions. Subscribe + with `solid_object room, payloads: :playmat_state` and handle the + `solid-objects:payload` DOM event. ERB component refreshes remain the default + and are unchanged. A mutation that changes payload state without changing a + declared observable still invalidates subscribers, through a revision-only + broadcast that carries no observable value to the browser. + ## 0.5.2 - 2026-08-09 - Read the database clock once per transaction instead of once per step, and diff --git a/Gemfile.lock b/Gemfile.lock index d54a357..adbb972 100644 --- a/Gemfile.lock +++ b/Gemfile.lock @@ -1,7 +1,7 @@ PATH remote: . specs: - solid_objects (0.5.2) + solid_objects (0.6.0) actioncable (>= 8.0) actionpack (>= 8.0) actionview (>= 8.0) @@ -373,7 +373,7 @@ CHECKSUMS rubocop-rails-omakase (1.1.0) sha256=2af73ac8ee5852de2919abbd2618af9c15c19b512c4cfc1f9a5d3b6ef009109d ruby-progressbar (1.13.0) sha256=80fc9c47a9b640d6834e0dc7b3c94c9df37f08cb072b7761e4a71e22cff29b33 securerandom (0.4.1) sha256=cc5193d414a4341b6e225f0cb4446aceca8e50d5e1888743fac16987638ea0b1 - solid_objects (0.5.2) + solid_objects (0.6.0) sqlite3 (2.9.5-aarch64-linux-gnu) sha256=78075b6337d3d182c6d2b4691049ed45cd220826160c9ea18946bf6a1de200dc sqlite3 (2.9.5-aarch64-linux-musl) sha256=18c801185deb4adc01ddb281e8f672a39e3d1729979ca91e39439cd3eac0402d sqlite3 (2.9.5-arm-linux-gnu) sha256=1bdfca0c7d63998c60b0f4a8e3c8df2d33800ccc4abd2d612eddbbbc92a4c48b diff --git a/app/assets/javascripts/solid_objects/state_payload.js b/app/assets/javascripts/solid_objects/state_payload.js new file mode 100644 index 0000000..9ec9470 --- /dev/null +++ b/app/assets/javascripts/solid_objects/state_payload.js @@ -0,0 +1,69 @@ +const deliveredRevisions = new Map() + +class SolidObjectsPayloadElement extends HTMLElement { + connectedCallback() { + if (this.dataset.started === "true") return + + this.dataset.started = "true" + this.deliver() + } + + deliver() { + try { + const name = this.dataset.name + const revision = revisionFor(this) + const scope = this.closest("[id]") + if (!name || !revision || !scope) return + + const key = `${scope.id}:${name}` + if (!newerRevision(revision, deliveredRevisions.get(key))) return + + const payload = JSON.parse(this.textContent) + deliveredRevisions.set(key, revision) + scope.dispatchEvent( + new CustomEvent("solid-objects:payload", { + bubbles: true, + detail: { + name, + instanceId: revision[0], + revision: revision[1], + payload + } + }) + ) + } catch { + this.dispatchEvent( + new CustomEvent("solid-objects:payload-error", { + bubbles: true, + detail: { reason: "invalid_payload" } + }) + ) + } finally { + this.remove() + } + } +} + +function newerRevision(candidate, current) { + if (!current) return true + + return candidate[0] > current[0] || + (candidate[0] === current[0] && candidate[1] > current[1]) +} + +function revisionFor(element) { + const revision = element.dataset.revision + if (!revision) return + + const values = revision.split(":").map(Number) + if ( + values.length !== 2 || + values.some((value) => !Number.isSafeInteger(value) || value < 0) + ) return + + return values +} + +if (!customElements.get("solid-objects-payload")) { + customElements.define("solid-objects-payload", SolidObjectsPayloadElement) +} diff --git a/app/helpers/solid_objects/actor_helper.rb b/app/helpers/solid_objects/actor_helper.rb index 11e3dff..fb2358e 100644 --- a/app/helpers/solid_objects/actor_helper.rb +++ b/app/helpers/solid_objects/actor_helper.rb @@ -2,8 +2,9 @@ module SolidObjects module ActorHelper - # @rbs (Reference, ?authorization_context: untyped) { (ActorView) -> untyped } -> untyped - def solid_object(reference, authorization_context: self, &block) + # @rbs (Reference, ?authorization_context: untyped, ?payloads: untyped) { (ActorView) -> untyped } -> untyped + def solid_object(reference, authorization_context: self, payloads: nil, &block) + payload_names = Array(payloads).map(&:to_s).uniq.presence actor = ActorView.new( reference:, view_context: self, @@ -13,7 +14,8 @@ def solid_object(reference, authorization_context: self, &block) subscription_data = { token: StreamToken.generate( reference, - observables: actor.scalar_observable_names + observables: actor.scalar_observable_names, + payloads: payload_names ) } if actor.component_tokens.any? @@ -30,10 +32,17 @@ def solid_object(reference, authorization_context: self, &block) data: { turbo_track: "reload" } ) end + payload_client = if payload_names + javascript_include_tag( + "solid_objects/state_payload", + type: "module", + data: { turbo_track: "reload" } + ) + end content_tag( :div, - safe_join([ refresh_client, subscription, content ].compact), + safe_join([ refresh_client, payload_client, subscription, content ].compact), id: DomIdentity.scope(reference) ) end diff --git a/docs/realtime.md b/docs/realtime.md index 29b9dac..a4d78b3 100644 --- a/docs/realtime.md +++ b/docs/realtime.md @@ -110,6 +110,93 @@ applications discover the namespaced engine asset. Applications created with explicitly serve the module. Turbo's normal morph rules still apply; use `data-turbo-permanent` for elements that must never be changed. +## Personalized state payloads + +Reactive ERB components cost one browser request per changed component. When a +single actor mutation changes several components, an application pays several +round trips for one logical update. A payload broadcast collapses that into one +message on the stream the page already has open. + +Declare the payload on the actor. The block receives the actor and the +subscriber's authorization context, and it runs **once per subscriber**, so two +sessions watching the same actor never see each other's private state: + +```ruby +class PlaymatRoom < SolidObjects::Actor + actor_type "playmat_room" + + attribute :hands, default: -> { {} } + attribute :turn, default: 1 + + observable :turn + + broadcast_payload :playmat_state do |room, authorization_context| + { + "turn" => room.turn, + "hand" => room.hands.fetch(authorization_context.session_id, []) + } + end +end +``` + +Subscribe the scope to it: + +```erb +<%= solid_object room, payloads: :playmat_state do |actor| %> +
+<% end %> +``` + +Handle it with any JavaScript. The gem dispatches a DOM event and requires no +framework: + +```javascript +document.addEventListener("solid-objects:payload", (event) => { + const { name, revision, payload } = event.detail + if (name !== "playmat_state") return + + renderPlaymat(payload) +}) +``` + +### What the protocol guarantees + +The payload travels as a Turbo Stream element on the existing actor stream, so +applications do not run a second WebSocket system. Each message carries the +actor identity plus the `instance_id` and monotonic `state_revision` that fence +component refreshes, and both the channel and the browser drop a payload that +is not newer than the last one delivered for that scope and name. A reconnecting +client receives the current payload on subscribe. + +Authorization is the same `authorize_query` boundary that components use, called +with the payload name and the subscriber's Cable connection. A subscriber that +fails the check is skipped rather than served a partial payload, and the payload +name is signed into the stream token, so a browser cannot ask for a payload the +server did not offer. + +Payload blocks read committed actor state through the same snapshot components +use. They cannot write application records, and the return value must be a JSON +object or array so the wire format stays inspectable. + +### mtg-playmat before and after + +Before, one mutation that touched three observables produced three refresh +elements and three HTTP requests: + +``` +commit -> 3 Action Cable messages -> 3 GET /solid_objects/components -> 3 renders +``` + +After, the same mutation delivers one personalized payload and the page renders +once: + +``` +commit -> 1 Action Cable message -> 0 HTTP requests -> 1 render +``` + +Components remain the default. An actor with no `broadcast_payload` and a scope +with no `payloads:` option behave exactly as before. + ## Authorization The HTML contains a signed actor identity token. Signing prevents modification; diff --git a/lib/solid_objects.rb b/lib/solid_objects.rb index db26a27..11bc94b 100644 --- a/lib/solid_objects.rb +++ b/lib/solid_objects.rb @@ -40,6 +40,7 @@ require "solid_objects/component_subscriptions" require "solid_objects/component_view" require "solid_objects/component_renderer" +require "solid_objects/payload_broadcast" require "solid_objects/state_snapshot" require "solid_objects/actor_view" require "solid_objects/actor_channel" diff --git a/lib/solid_objects/actor.rb b/lib/solid_objects/actor.rb index 7ccc015..ba5f788 100644 --- a/lib/solid_objects/actor.rb +++ b/lib/solid_objects/actor.rb @@ -57,6 +57,13 @@ def observable(name, &block) definition.add_observable(name, block) end + # @rbs (Symbol | String) { (untyped, untyped) -> untyped } -> ActorDefinition::Handler + def broadcast_payload(name, &block) + raise InvalidActor, "payload broadcasts require a block" unless block + + definition.add_payload_broadcast(name, block) + end + # @rbs (?Integer) -> Integer def state_version(version = nil) definition.set_state_version(version) if version diff --git a/lib/solid_objects/actor_channel.rb b/lib/solid_objects/actor_channel.rb index cb15640..d91873f 100644 --- a/lib/solid_objects/actor_channel.rb +++ b/lib/solid_objects/actor_channel.rb @@ -19,7 +19,9 @@ def subscribed @reference = Reference.new(actor_type:, actor_id:) @scalar_observables = identity["observables"] + @payload_names = identity["payloads"] validate_scalar_observables! + validate_payload_names! @component_subscriptions = ComponentSubscriptions.parse( params["components"], reference: @@ -36,6 +38,7 @@ def subscribed ) end refresh_outdated_components(snapshot) + transmit_state_payloads(snapshot) rescue KeyError, JSON::ParserError, InvalidStreamToken, @@ -46,14 +49,20 @@ def subscribed private - attr_reader :reference, :component_subscriptions, :scalar_observables + attr_reader :reference, + :component_subscriptions, + :scalar_observables, + :payload_names # @rbs (String) -> void def receive_broadcast(stream) invalidation = TurboStreamRenderer.invalidation(stream) - if !invalidation || - scalar_observables.nil? || - scalar_observables.include?(invalidation.fetch("observable_name")) + revision_only = invalidation && + invalidation.fetch("observable_name") == PayloadBroadcast::REVISION_OBSERVABLE + if !revision_only && + (!invalidation || + scalar_observables.nil? || + scalar_observables.include?(invalidation.fetch("observable_name"))) transmit stream end return unless invalidation @@ -61,6 +70,48 @@ def receive_broadcast(stream) component_subscriptions .refreshes_for(invalidation) .each { |refresh| transmit refresh } + transmit_state_payloads(ActorSnapshot.new(reference)) + end + + # @rbs (ActorSnapshot) -> void + def transmit_state_payloads(snapshot) + return if payload_names.nil? || payload_names.empty? + return unless newer_payload_revision?(snapshot) + + payload_names.each do |name| + payload = PayloadBroadcast.new( + snapshot:, + name:, + authorization_context: connection + ).call + transmit TurboStreamRenderer.state_payload(payload) + rescue Unauthorized + next + end + @payload_revision = [ snapshot.instance_id, snapshot.revision ] + end + + # @rbs (ActorSnapshot) -> bool + def newer_payload_revision?(snapshot) + current = @payload_revision + return true unless current + + (current <=> [ snapshot.instance_id, snapshot.revision ]) == -1 + end + + # @rbs () -> void + def validate_payload_names! + return unless payload_names + + broadcasts = SolidObjects + .registry + .fetch(reference.actor_type) + .definition + .payload_broadcasts + unknown = payload_names.find { |name| !broadcasts.key?(name.to_sym) } + return unless unknown + + raise InvalidStreamToken, "unknown payload broadcast #{unknown.inspect}" end # @rbs (ActorSnapshot) -> void diff --git a/lib/solid_objects/actor_definition.rb b/lib/solid_objects/actor_definition.rb index 3d868ca..46aa9c4 100644 --- a/lib/solid_objects/actor_definition.rb +++ b/lib/solid_objects/actor_definition.rb @@ -9,6 +9,7 @@ class ActorDefinition # @rbs @messages: Hash[Symbol, Handler] # @rbs @queries: Hash[Symbol, Handler] # @rbs @observables: Hash[Symbol, Handler] + # @rbs @payload_broadcasts: Hash[Symbol, Handler] # @rbs @state_version: Integer # @rbs @state_migrations: Array[StateMigration] # @rbs @activation_hooks: Array[Proc] @@ -20,6 +21,7 @@ class ActorDefinition :messages, :queries, :observables, + :payload_broadcasts, :state_version, :state_migrations, :activation_hooks, @@ -31,6 +33,7 @@ def initialize @messages = {} @queries = {} @observables = {} + @payload_broadcasts = {} @state_version = 1 @state_migrations = [] @activation_hooks = [] @@ -83,6 +86,21 @@ def add_observable(name, block = nil) end end + # @rbs (Symbol | String, Proc) -> Handler + def add_payload_broadcast(name, block) + payload_name = name.to_sym + if payload_broadcasts.key?(payload_name) + raise InvalidActor, "#{payload_name.inspect} payload broadcast is already defined" + end + unless payload_name.to_s.match?(/\A[a-zA-Z0-9_]+\z/) + raise InvalidActor, "payload broadcast names may contain only letters, digits, and underscores" + end + + Handler.new(name: payload_name, block:).tap do |handler| + payload_broadcasts[payload_name] = handler + end + end + # @rbs (Class) -> ActorDefinition def synchronize_instance_messages(actor_class) names = actor_message_method_names(actor_class) @@ -149,6 +167,7 @@ def duplicate copy.instance_variable_set(:@messages, messages.dup) copy.instance_variable_set(:@queries, queries.dup) copy.instance_variable_set(:@observables, observables.dup) + copy.instance_variable_set(:@payload_broadcasts, payload_broadcasts.dup) copy.instance_variable_set(:@state_version, state_version) copy.instance_variable_set(:@state_migrations, state_migrations.dup) copy.instance_variable_set(:@activation_hooks, activation_hooks.dup) diff --git a/lib/solid_objects/errors.rb b/lib/solid_objects/errors.rb index 73b40d9..f4c63c5 100644 --- a/lib/solid_objects/errors.rb +++ b/lib/solid_objects/errors.rb @@ -43,6 +43,12 @@ class InvalidComponentToken < Error class UnknownComponent < Error end + class UnknownPayloadBroadcast < Error + end + + class InvalidPayloadBroadcast < Error + end + class UnknownComponentDependency < Error end diff --git a/lib/solid_objects/executor.rb b/lib/solid_objects/executor.rb index ce3dbb9..ba90e13 100644 --- a/lib/solid_objects/executor.rb +++ b/lib/solid_objects/executor.rb @@ -27,7 +27,7 @@ def call result = invoke_actor(message_context) ensure_query_did_not_mutate_state!(state_before) observable_changes = changed_observables(observables_before, actor.observable_values) - complete(result, observable_changes) + complete(result, observable_changes, state_changed: actor.state.to_h != state_before) true rescue LostActivation raise @@ -75,8 +75,8 @@ def changed_observables(before, after) end end - # @rbs (untyped, Hash[String, untyped]) -> void - def complete(result, observable_changes) + # @rbs (untyped, Hash[String, untyped], state_changed: bool) -> void + def complete(result, observable_changes, state_changed:) serialized_state = Serialization.dump( actor.state.to_h, max_bytes: SolidObjects.configuration.max_state_bytes @@ -111,7 +111,12 @@ def complete(result, observable_changes) enqueued_effects.concat( enqueue_actor_messages(locked_message, instance, outbound_message_intents) ) - enqueue_broadcasts(locked_message, instance, observable_changes) + enqueue_broadcasts( + locked_message, + instance, + observable_changes, + state_changed: + ) claimed_message.destroy! end @@ -249,9 +254,14 @@ def enqueue_actor_messages(locked_message, instance, intents) end end - # @rbs (Message, Instance, Hash[String, untyped]) -> void - def enqueue_broadcasts(locked_message, instance, observable_changes) - observable_changes.each do |observable_name, value| + # @rbs (Message, Instance, Hash[String, untyped], state_changed: bool) -> void + def enqueue_broadcasts(locked_message, instance, observable_changes, state_changed:) + broadcasts = observable_changes + if broadcasts.empty? && state_changed && payload_broadcasts? + broadcasts = { PayloadBroadcast::REVISION_OBSERVABLE => {} } + end + + broadcasts.each do |observable_name, value| Broadcast.create!( message: locked_message, instance:, @@ -266,6 +276,11 @@ def enqueue_broadcasts(locked_message, instance, observable_changes) end end + # @rbs () -> bool + def payload_broadcasts? + actor.class.definition.payload_broadcasts.any? + end + # @rbs (Exception) -> void def fail_message(error) error_details = serialized_error(error) diff --git a/lib/solid_objects/payload_broadcast.rb b/lib/solid_objects/payload_broadcast.rb new file mode 100644 index 0000000..d62af38 --- /dev/null +++ b/lib/solid_objects/payload_broadcast.rb @@ -0,0 +1,67 @@ +# rbs_inline: enabled + +module SolidObjects + class PayloadBroadcast + MAXIMUM_PAYLOAD_BYTES = 1_048_576 + REVISION_OBSERVABLE = "solid_objects.revision" + + # @rbs @snapshot: ActorSnapshot + # @rbs @name: String + # @rbs @authorization_context: untyped + + attr_reader :name + + # @rbs (snapshot: ActorSnapshot, name: String, authorization_context: untyped) -> void + def initialize(snapshot:, name:, authorization_context:) + @snapshot = snapshot + @name = name + @authorization_context = authorization_context + end + + # @rbs () -> Hash[String, untyped] + def call + handler = snapshot.actor_class.definition.payload_broadcasts[name.to_sym] + raise UnknownPayloadBroadcast, "unknown payload broadcast #{name.inspect}" unless handler + + authorize! + { + "actor_type" => snapshot.reference.actor_type, + "actor_id" => snapshot.reference.actor_id, + "name" => name, + "instance_id" => snapshot.instance_id, + "revision" => snapshot.revision, + "payload" => rendered_payload(handler) + } + end + + private + + attr_reader :snapshot, :authorization_context + + # @rbs (ActorDefinition::Handler) -> untyped + def rendered_payload(handler) + payload = Serialization.dump( + handler.block.call(snapshot.actor, authorization_context), + max_bytes: MAXIMUM_PAYLOAD_BYTES + ) + return payload if payload.is_a?(Hash) || payload.is_a?(Array) + + raise InvalidPayloadBroadcast, + "payload broadcast #{name.inspect} must return a JSON object or array" + end + + # @rbs () -> void + def authorize! + authorized = SolidObjects.configuration.authorize_query.call( + actor_type: snapshot.reference.actor_type, + actor_id: snapshot.reference.actor_id, + message_name: name, + arguments: {}, + authorization_context: + ) + return if authorized + + raise Unauthorized, "actor payload broadcast is not authorized" + end + end +end diff --git a/lib/solid_objects/stream_token.rb b/lib/solid_objects/stream_token.rb index f86ae00..cb8e7c3 100644 --- a/lib/solid_objects/stream_token.rb +++ b/lib/solid_objects/stream_token.rb @@ -9,13 +9,14 @@ module StreamToken module_function - # @rbs (Reference, ?observables: Array[String]?) -> String - def generate(reference, observables: nil) + # @rbs (Reference, ?observables: Array[String]?, ?payloads: Array[String]?) -> String + def generate(reference, observables: nil, payloads: nil) identity = { "actor_type" => reference.actor_type, "actor_id" => reference.actor_id } identity["observables"] = observables if observables + identity["payloads"] = payloads if payloads validate_identity!(identity) verifier.generate(identity, purpose: PURPOSE) end @@ -36,21 +37,27 @@ def validate_identity!(identity) raise InvalidStreamToken, "invalid actor stream token" end - observables = identity["observables"] - return identity unless observables + validate_names!(identity["observables"], "observables") + validate_names!(identity["payloads"], "payloads") + identity + end + + # @rbs (untyped, String) -> void + def validate_names!(names, label) + return unless names - valid = observables.is_a?(Array) && - observables.length <= MAXIMUM_OBSERVABLES && - observables.uniq.length == observables.length && - observables.all? do |observable| - observable.is_a?(String) && - observable.match?(/\A[a-zA-Z0-9_]+\z/) + valid = names.is_a?(Array) && + names.length <= MAXIMUM_OBSERVABLES && + names.uniq.length == names.length && + names.all? do |name| + name.is_a?(String) && name.match?(/\A[a-zA-Z0-9_]+\z/) end - return identity if valid + return if valid - raise InvalidStreamToken, "invalid actor stream observables" + raise InvalidStreamToken, "invalid actor stream #{label}" end private_class_method :validate_identity! + private_class_method :validate_names! # @rbs () -> ActiveSupport::MessageVerifier def verifier diff --git a/lib/solid_objects/turbo_stream_renderer.rb b/lib/solid_objects/turbo_stream_renderer.rb index a03463a..d4b1c96 100644 --- a/lib/solid_objects/turbo_stream_renderer.rb +++ b/lib/solid_objects/turbo_stream_renderer.rb @@ -16,11 +16,15 @@ def observable(broadcast) actor_type: broadcast.instance.actor_type, actor_id: broadcast.instance.actor_id ) - stream = observable_value( - reference, - broadcast.observable_name, - broadcast.value - ) + stream = if broadcast.observable_name == PayloadBroadcast::REVISION_OBSERVABLE + "" + else + observable_value( + reference, + broadcast.observable_name, + broadcast.value + ) + end metadata = Base64.urlsafe_encode64( JSON.generate( "instance_id" => broadcast.instance_id, @@ -50,6 +54,19 @@ def component_refresh(registration, instance_id, revision) %(