diff --git a/src/off_topic/internal/runtime.gleam b/src/off_topic/internal/runtime.gleam index f1c11b3..07b7cc0 100644 --- a/src/off_topic/internal/runtime.gleam +++ b/src/off_topic/internal/runtime.gleam @@ -65,9 +65,13 @@ pub type Dependency /// Use with [init](#init), [update](#update), and [view](#view) when wiring /// up off_topic manually instead of through [application](#application). pub opaque type Model(model, message) { - Dormant(model: model) - Application(model: model, tick: Int, state: Running(message)) - Component(model: model, tick: Int, state: Running(message), clients: Int) + Model( + tick: Int, + app: model, + mounted: Bool, + state: Running(message), + connected_clients: Int, + ) } /// The subscription runtime's wrapper around your application message type. @@ -76,11 +80,21 @@ pub opaque type Model(model, message) { /// up off_topic manually instead of through [application](#application). pub opaque type Message(message) { ClientRuntimeMounted - ToApp(message: message) - GotStarted(tick: Int, path: List(Int), cleanup: fn() -> Nil) - FromRemote(tick: Int, path: List(Int), message: Dynamic) - ThrottleDelay(tick: Int, path: List(Int), message: message, ts: Int) - DelayFired(tick: Int, path: List(Int), delay_tick: Int, message: message) + GotAppMessage(message: message) + SubscriptionStarted(tick: Int, path: List(Int), cleanup: fn() -> Nil) + RemoteSubscriptionSentMessage(tick: Int, path: List(Int), message: Dynamic) + SubscriptionSentDeferredMessage( + tick: Int, + path: List(Int), + message: message, + ts: Int, + ) + DelayTimeoutElapsed( + tick: Int, + path: List(Int), + delay_tick: Int, + message: message, + ) ComponentConnected(forward: Option(message)) ComponentDisconnected(forward: Option(message)) } @@ -326,24 +340,8 @@ pub fn application( subscriptions app_subscriptions: fn(model) -> Subscription(message), view app_view: fn(model) -> Element(message), ) -> App(flags, Model(model, message), Message(message)) { - let #(update, view) = wrap_runtime(app_update, app_subscriptions, app_view) - - // Applications always start subscriptions immediately — no shadow root check. - let init = fn(flags) { - let #(app_model, app_effect) = app_init(flags) - case lustre.is_browser() { - True -> { - let eff = - effect.batch([ - effect.map(app_effect, ToApp), - effect.from(fn(dispatch) { dispatch(ClientRuntimeMounted) }), - ]) - #(Dormant(model: app_model), eff) - } - False -> #(Dormant(model: app_model), effect.map(app_effect, ToApp)) - } - } - + let #(init, update, view) = + wrap_runtime(app_init, app_update, app_subscriptions, app_view) lustre.application(init:, update:, view:) } @@ -362,13 +360,8 @@ pub fn component( view app_view: fn(model) -> Element(message), options options: List(component.Option(message)), ) -> App(Nil, Model(model, message), Message(message)) { - let #(update, view) = wrap_runtime(app_update, app_subscriptions, app_view) - - // Components never dispatch Mounted — lifecycle is owned by ComponentConnected. - let init = fn(flags) { - let #(app_model, app_effect) = app_init(flags) - #(Dormant(model: app_model), effect.map(app_effect, ToApp)) - } + let #(init, update, view) = + wrap_runtime(app_init, app_update, app_subscriptions, app_view) let #(user_on_connect, user_on_disconnect) = component.extract_lifecycle(options) @@ -376,23 +369,24 @@ pub fn component( // User options lifted to Message(message) via ToApp, followed by our runtime // lifecycle options. Lustre's configure is last-wins, so appending last means // our on_connect / on_disconnect always take effect. - let lustre_opts = - list.map(options, component.map(_, ToApp)) + let options = + list.map(options, component.map(_, GotAppMessage)) |> list.append([ component.on_connect(ComponentConnected(user_on_connect)), component.on_disconnect(ComponentDisconnected(user_on_disconnect)), ]) |> component.to_lustre_options - lustre.component(init:, update:, view:, options: lustre_opts) + lustre.component(init:, update:, view:, options:) } -fn wrap_runtime(app_update, app_subscriptions, app_view) { +fn wrap_runtime(app_init, app_update, app_subscriptions, app_view) { + let init = init(_, app_init) let view = view(_, app_view) let update = fn(model, msg) { update(model, msg, app_subscriptions, app_update) } - #(update, view) + #(init, update, view) } // -- APPLICATION BUILDING BLOCKS --------------------------------------------- @@ -409,28 +403,31 @@ pub fn init( ) -> #(Model(model, message), Effect(Message(message))) { let #(model, app_effect) = init(flags) - case lustre.is_browser() { - True -> { - // Detect at runtime whether we're in a component (shadow root) or a plain - // application. Components defer mounting to ComponentConnected; plain - // applications dispatch Mounted immediately so subscriptions start now. - let check_mount = - effect.before_paint(fn(dispatch, root) { - case is_shadow_root(root) { - False -> dispatch(ClientRuntimeMounted) - True -> Nil - } - }) - let effect = effect.batch([effect.map(app_effect, ToApp), check_mount]) - #(Dormant(model:), effect) - } - False -> #(Dormant(model:), effect.map(app_effect, ToApp)) - } -} + // Figure out if we will receive connect/disconnect messages. + // + // We will always receive a ClientRuntimeMounted mesasge: + // - Server components don't fire before_paint, but will emit an event + // - For applications and components, before_paint always fires synchronously, + // but AFTER a potential connected message, if we do receive those. + // + // Waiting until the Mounted message fired, we can then check if we just received a + // connected message prior and if we did, we are in component mode! + let effect = + effect.batch([ + effect.map(app_effect, GotAppMessage), + effect.before_paint(fn(dispatch, _root) { dispatch(ClientRuntimeMounted) }), + ]) -@external(javascript, "../../off_topic_ffi.mjs", "is_shadow_root") -fn is_shadow_root(_root: Dynamic) -> Bool { - False + let model = + Model( + tick: 0, + app: model, + mounted: False, + state: Group([]), + connected_clients: 0, + ) + + #(model, effect) } /// Handle a message in the subscription runtime. @@ -445,90 +442,75 @@ pub fn update( subscriptions: fn(model) -> Subscription(message), update: fn(model, message) -> #(model, Effect(message)), ) -> #(Model(model, message), Effect(Message(message))) { - case message { - ToApp(message:) -> to_app(model, message, subscriptions, update) + case echo message { + GotAppMessage(message:) -> + to_app(model, model.state, message, subscriptions, update) - ClientRuntimeMounted -> + // client components receive mounted exactly once on startup, + // server components can repeatedly receive mounted whenever a client joins (and have to ignore connected instead) + ClientRuntimeMounted -> { case model { - Dormant(model: app_model) -> { - let #(state, effects) = start_subscriptions(app_model, subscriptions) - #(Application(model: app_model, tick: 0, state:), effects) + // Server component — a new client joined while subscriptions are live. + // Re-emit remote subscriptions so the new client receives them. + Model(mounted: True, connected_clients:, ..) if connected_clients > 0 -> { + #(model, effect.batch(re_emit(model.state, [], []))) } - _ -> #(model, effect.none()) - } - ComponentConnected(forward:) -> - case model { - Dormant(model: app_model) -> { - // Apply the user's connect message first so the initial subscription - // set reflects the post-connect model, avoiding a mount-then-diff. - let #(final_model, user_effect) = case forward { - None -> #(app_model, effect.none()) - Some(msg) -> { - let #(m, e) = update(app_model, msg) - #(m, effect.map(e, ToApp)) - } - } - let #(state, mount_effects) = - start_subscriptions(final_model, subscriptions) - let model = Component(model: final_model, tick: 0, state:, clients: 1) - #(model, effect.batch([user_effect, mount_effects])) - } - Application(..) -> - // Application-started — subscriptions are permanent, just forward. - maybe_forward(model, forward, subscriptions, update) - Component(clients:, state:, ..) -> { - let re_emit_effects = effect.batch(re_emit(state, [], [])) - let model = Component(..model, clients: clients + 1) - let #(model, forward_effect) = - maybe_forward(model, forward, subscriptions, update) - #(model, effect.batch([re_emit_effects, forward_effect])) + // All other cases: + // + // (False, 0) Plain application — no component lifecycle, subscriptions run forever. + // + // (False, _) Custom element or server component — clients connected before mount, + // so subscriptions weren't started yet. Start them now. + // + // (True, 0) Server component re-mounted after all clients disconnected. + // Subscriptions were torn down on last disconnect; restart them. + _ -> { + let connected_clients = int.max(1, model.connected_clients) + step_lifecycle(model, True, connected_clients, subscriptions) } } + } - ComponentDisconnected(forward:) -> - case model { - Dormant(..) | Application(..) -> - // Dormant: no subscriptions to stop; Application: they run forever. - maybe_forward(model, forward, subscriptions, update) - Component(model: app_model, state:, clients: 1, ..) -> { - // Last client disconnected — stop subscriptions then forward. - let cleanup_effect = effect.batch(cleanup(state, [], [])) - let #(dormant, forward_effect) = - maybe_forward(Dormant(app_model), forward, subscriptions, update) - #(dormant, effect.batch([cleanup_effect, forward_effect])) - } - Component(clients:, ..) -> { - let state = Component(..model, clients: clients - 1) - maybe_forward(state, forward, subscriptions, update) - } - } + ComponentConnected(msg) -> { + // Client component going from 0→1 while already mounted — restart subscriptions. + // Apply forward first so the initial subscription set reflects the post-connect model. + // + // Otherwise subscriptions are running, will start on mount, or this is a server component. + let connected_clients = model.connected_clients + 1 + // When we hit this in a browser while not mounted , we know for sure that + // we are in a custom element, so we can mount early as a slight optimisation. + // The mounted dispatch afterwards will re_emit, which we ignore. + let mounted = model.mounted || lustre.is_browser() + forward(msg, model, mounted, connected_clients, subscriptions, update) + } - ThrottleDelay(tick:, path:, message:, ts:) -> { - case - get_state(model) - |> result.try(throttle_delay(ts, _, tick, path, message)) - { + ComponentDisconnected(msg) -> { + let mounted = model.mounted + let connected_clients = int.max(0, model.connected_clients - 1) + forward(msg, model, mounted, connected_clients, subscriptions, update) + } + + SubscriptionSentDeferredMessage(tick:, path:, message:, ts:) -> { + case deferred(ts, tick, path, model.state, message) { Ok(#(state, Immediate(message:))) -> - to_app(set_state(model, state), message, subscriptions, update) - Ok(#(state, Delayed(message:, delay:, delay_tick:))) -> #( - set_state(model, state), - after(delay, DelayFired(tick:, path:, delay_tick:, message:)), - ) + to_app(model, state, message, subscriptions, update) + Ok(#(state, Delayed(message:, delay:, delay_tick:))) -> { + let delay_message = + DelayTimeoutElapsed(tick:, path:, delay_tick:, message:) + #(Model(..model, state:), after(delay, delay_message)) + } Error(Nil) -> #(model, effect.none()) } } - GotStarted(tick:, path:, cleanup:) -> { - let new_model = - get_state(model) - |> result.try(started(tick, path, _, cleanup)) - |> result.map(set_state(model, _)) - |> result.unwrap(model) - #(new_model, effect.none()) - } + SubscriptionStarted(tick:, path:, cleanup:) -> + case started(tick, path, model.state, cleanup) { + Ok(state) -> #(Model(..model, state:), effect.none()) + Error(_) -> #(model, effect.from(fn(_) { cleanup() })) + } - FromRemote(tick:, path:, message:) -> { + RemoteSubscriptionSentMessage(tick:, path:, message:) -> { use state <- emit_message(model, path, subscriptions, update) case state { Listening(decoder:, ..) if state.tick == tick -> @@ -537,7 +519,7 @@ pub fn update( } } - DelayFired(tick:, path:, delay_tick:, message:) -> { + DelayTimeoutElapsed(tick:, path:, delay_tick:, message:) -> { use state <- emit_message(model, path, subscriptions, update) case state { Running(..) if state.tick == tick && state.delay_tick == delay_tick -> @@ -548,63 +530,6 @@ pub fn update( } } -fn emit_message(model, path, subscriptions, update, f) { - case - get_state(model) - |> result.try(read_at(_, path, f)) - { - Ok(msg) -> to_app(model, msg, subscriptions, update) - Error(_) -> #(model, effect.none()) - } -} - -fn read_at(state, path, f) { - case path, state { - [], _ -> f(state) - [index, ..rest], Group(children) -> - case at(children, index) { - Ok(child) -> read_at(child, rest, f) - Error(Nil) -> Error(Nil) - } - _, _ -> Error(Nil) - } -} - -fn at(list, index) { - case list { - [] -> Error(Nil) - [first, ..] if index <= 0 -> Ok(first) - [_, ..rest] -> at(rest, index - 1) - } -} - -fn start_subscriptions(app_model, subscriptions) { - let #(state, effects) = start(0, [], subscriptions(app_model), []) - #(state, effect.batch(effects)) -} - -fn get_state(model: Model(a, b)) -> Result(Running(b), Nil) { - case model { - Application(state:, ..) | Component(state:, ..) -> Ok(state) - Dormant(..) -> Error(Nil) - } -} - -fn set_state(model: Model(a, b), state: Running(b)) -> Model(a, b) { - case model { - Application(..) -> Application(..model, state:) - Component(..) -> Component(..model, state:) - Dormant(..) -> model - } -} - -fn maybe_forward(model, forward, subscriptions, update) { - case forward { - None -> #(model, effect.none()) - Some(msg) -> to_app(model, msg, subscriptions, update) - } -} - /// Render the application view through the subscription runtime. /// /// The low-level counterpart to [application](#application), for use when you @@ -624,7 +549,7 @@ pub fn view( use tick <- decode.field("tick", decode.int) use path <- decode.field("path", decode.list(decode.int)) use message <- decode.field("message", decode.dynamic) - decode.success(FromRemote(tick:, path:, message:)) + decode.success(RemoteSubscriptionSentMessage(tick:, path:, message:)) } let on_dispatch = @@ -638,70 +563,110 @@ pub fn view( } } - let app = case model { - Dormant(model:) -> view(model) - Application(model:, tick:, ..) | Component(model:, tick:, ..) -> { - // we increment the `tick` whenever we update the subscriptions, which - // happens whenever the user model got updated. - // This means we can skip "internal" updates in our runtime, which never - // influence the view, allowing us to not re-render throttled events etc. - use <- element.memo([element.ref([tick])]) - view(model) - } + let app = { + // we increment the `tick` whenever we update the subscriptions, which + // happens whenever the user model got updated. + // This means we can skip "internal" updates in our runtime, which never + // influence the view, allowing us to not re-render throttled events etc. + element.memo([element.ref([model.tick])], fn() { view(model.app) }) } - element.fragment([runtime, element.map(app, ToApp)]) + element.fragment([runtime, element.map(app, GotAppMessage)]) } // -- SUB DISPATCH ------------------------------------------------------------ -fn to_app(model: Model(_, _), message, subscriptions, update) { - let #(app_model, app_effect) = update(model.model, message) +fn to_app(model: Model(_, _), state, message, subscriptions, update) { + let Model(tick:, app:, mounted:, connected_clients:, ..) = model + let #(app, app_effect) = update(app, message) + let #(model, sync_effect) = + step(tick, app, mounted, state, connected_clients, subscriptions) + #(model, effect(sync_effect, app_effect)) +} + +fn effect(sync_effect, app_effect) -> Effect(Message(message)) { + effect.batch([sync_effect, effect.map(app_effect, GotAppMessage)]) +} - case model { - Dormant(..) -> #(Dormant(app_model), effect.map(app_effect, ToApp)) - Application(state:, tick:, ..) | Component(state:, tick:, ..) -> - to_app_and_tick(tick, state, subscriptions, app_model, app_effect, model) +fn step_lifecycle(model, mounted, connected_clients, subscriptions) { + let Model(tick:, app:, state:, ..) = model + step(tick, app, mounted, state, connected_clients, subscriptions) +} + +fn forward(forward, model, mounted, connected_clients, subscriptions, update) { + case forward { + Some(message) -> { + let Model(tick:, state:, app: app_model, ..) = model + let #(app_model, app_effect) = update(app_model, message) + let #(model, sync_effect) = + step(tick, app_model, mounted, state, connected_clients, subscriptions) + #(model, effect(sync_effect, app_effect)) + } + None -> step_lifecycle(model, mounted, connected_clients, subscriptions) } } -fn to_app_and_tick(tick, state, subscriptions, app_model, app_effect, model) { - let tick = tick + 1 - let #(state, sync_effect) = - diff(tick, [], state, subscriptions(app_model), []) - let effect = effect.batch([effect.map(app_effect, ToApp), ..sync_effect]) - let model = case model { - Application(..) -> Application(model: app_model, tick:, state:) - Component(clients:, ..) -> - Component(model: app_model, tick:, state:, clients:) - Dormant(..) -> Dormant(app_model) +fn emit_message(model: Model(_, _), path, subscriptions, update, f) { + let state = model.state + case result.try(get(state, path), f) { + Ok(msg) -> to_app(model, state, msg, subscriptions, update) + Error(_) -> #(model, effect.none()) } - #(model, effect) } -fn update_at(state, path, f) { +fn get(state, path) { case path, state { - [], _ -> f(state) - [index, ..rest], Group(children) -> { - use #(prefix, child, suffix) <- result.try(split(children, index, [])) - use #(child, value) <- result.try(update_at(child, rest, f)) - Ok(#(Group(unsplit(prefix, child, suffix)), value)) - } + [], _ -> Ok(state) + [index, ..rest], Group(children) -> + result.try(at(children, index), get(_, rest)) _, _ -> Error(Nil) } } +fn at(list, index) { + case list { + [] -> Error(Nil) + [first, ..] if index <= 0 -> Ok(first) + [_, ..rest] -> at(rest, index - 1) + } +} + +fn step(tick, app, mounted, state, connected_clients, subscriptions) { + // let model = Model(..model, mounted:, connected_clients:) + // step(model, model.app, subscriptions) + + // if we are not mounted or have no connected clients, we treat everything + // as if there are no subscriptions - + // this will automatically re-start or stop the current state accordingly + // as clients connect and disconnect. + let subscriptions = case mounted && connected_clients > 0 { + True -> subscriptions(app) + False -> none() + } + + let tick = next_tick(tick) + let #(state, effects) = diff(tick, [], state, subscriptions, []) + let model = Model(tick:, state:, app:, mounted:, connected_clients:) + #(model, effect.batch(effects)) +} + +fn next_tick(tick) { + case tick < 0xffffffff { + True -> tick + 1 + False -> 0 + } +} + fn started(tick, path, state, cleanup) { - use #(state, _) <- result.map( - update_at(state, path, fn(s) { - case s { - Running(..) if s.tick == tick -> - Ok(#(Running(..s, cleanup: cleanup), Nil)) - _ -> Error(Nil) - } - }), - ) - state + case path, state { + [], Running(..) if state.tick == tick -> Ok(Running(..state, cleanup:)) + [index, ..rest], Group(children) -> { + use #(prefix, child, suffix) <- result.try(split(children, index, [])) + use child <- result.try(started(tick, rest, child, cleanup)) + Ok(Group(unsplit(prefix, child, suffix))) + } + _, _ -> Error(Nil) + } } fn split(list, index, prefix) { @@ -721,28 +686,35 @@ type Dispatch(message) { Delayed(message: message, delay: Int, delay_tick: Int) } -fn throttle_delay(now, state, tick, path, message) { - use state <- update_at(state, path) - case state { - Running(config:, last_fired:, delay_tick:, ..) if state.tick == tick -> { - let throttle_passes = - config.throttle == 0 || now - last_fired >= config.throttle - let delay_tick = delay_tick + 1 - case throttle_passes, config.delay { - True, _ -> - Ok(#( - Running(..state, last_fired: now, delay_tick:), - Immediate(message), - )) - False, wait if wait > 0 -> - Ok(#( - Running(..state, delay_tick:), - Delayed(message, wait, delay_tick), - )) - False, _ -> Error(Nil) +fn deferred(now, tick, path, state, message) { + case path, state { + [], Running(config:, last_fired:, delay_tick:, ..) if state.tick == tick -> { + let Config(throttle:, delay:, ..) = config + + let delay_tick = next_tick(delay_tick) + let throttle_passes = throttle == 0 || now - last_fired >= throttle + + case throttle_passes, delay > 0 { + True, _ -> { + let state = Running(..state, last_fired: now, delay_tick:) + Ok(#(state, Immediate(message))) + } + False, True -> { + let state = Running(..state, delay_tick:) + Ok(#(state, Delayed(message, delay, delay_tick))) + } + False, False -> Error(Nil) } } - _ -> Error(Nil) + + [index, ..rest], Group(children) -> { + use #(prefix, child, suffix) <- result.try(split(children, index, [])) + use #(child, result) <- result.try({ + deferred(now, tick, rest, child, message) + }) + Ok(#(Group(unsplit(prefix, child, suffix)), result)) + } + _, _ -> Error(Nil) } } @@ -828,7 +800,7 @@ fn start(tick, path, subscription, effects) { effect.from(fn(dispatch) { let path = list.reverse(path) let cleanup = start(wrap_dispatch(dispatch, config, tick, path)) - dispatch(GotStarted(tick:, path:, cleanup:)) + dispatch(SubscriptionStarted(tick:, path:, cleanup:)) }) #(running(tick, False, config), [effect, ..effects]) } @@ -837,7 +809,7 @@ fn start(tick, path, subscription, effects) { effect.before_paint(fn(dispatch, root) { let path = list.reverse(path) let cleanup = start(wrap_dispatch(dispatch, config, tick, path), root) - dispatch(GotStarted(tick:, path:, cleanup:)) + dispatch(SubscriptionStarted(tick:, path:, cleanup:)) }) #(running(tick, True, config), [effect, ..effects]) } @@ -875,11 +847,11 @@ fn running(tick: Int, element: Bool, config: Config) -> Running(message) { fn wrap_dispatch(dispatch, config, tick, path) { case config { Config(throttle: 0, delay: 0, ..) -> fn(message) { - dispatch(ToApp(message)) + dispatch(GotAppMessage(message)) } _ -> fn(message) { let ts = now_ms() - dispatch(ThrottleDelay(tick:, path:, message:, ts:)) + dispatch(SubscriptionSentDeferredMessage(tick:, path:, message:, ts:)) } } } diff --git a/src/off_topic_ffi.mjs b/src/off_topic_ffi.mjs index 7c3ec65..37e57b4 100644 --- a/src/off_topic_ffi.mjs +++ b/src/off_topic_ffi.mjs @@ -2,10 +2,6 @@ export function send_after(timeout, message, dispatch) { setTimeout(dispatch, timeout, message); } -export function is_shadow_root(root) { - return root instanceof ShadowRoot; -} - export function now_ms() { return Date.now(); }