A server-side pub/sub system that lets external services (workers, webhooks, other actions) push updates to clients watching a specific entity. Built on top of the existing slot/SSE pipeline — no new transport.
A topic (Topic<TPayload>) loads the entity from the database when
triggered. State modifiers (StateModifier<TPayload, TState>) receive the
payload + the watcher's current slot state + a live lookup of all the client's
held states, and decide whether to push an update (return new state), skip
(return current), or revoke (return null).
Topic<TPayload> — loads the entity once per triggerpublic abstract class Topic<TPayload> : Object {
/** Load the payload for `topic_key` (e.g. "cat:42" → Cat from the database). */
public abstract async TPayload load(string topic_key) throws GLib.Error;
}
Runs once per trigger (not per watcher). DI-constructed — can inject databases, APIs, etc.
StateModifier<TPayload, TState> — derives slot state per watcherpublic abstract class StateModifier<TPayload, TState> : Object {
/**
* @param payload The trigger payload (loaded once by Topic<T>).
* @param current The current state of the slot being modified, or null if
* it has expired from the cache.
* @param held Live lookup of the client's currently-held states by
* type_name (reads CURRENT cache values at trigger time).
* @return New TState → push; `current` (same ref) → skip; null → clear + remove.
*/
public abstract TState? derive(TPayload payload, TState? current, HeldStates held);
}
HeldStates — live lookup of client state at trigger timepublic class HeldStates : Object {
/** Typed current state for a held type, or null if expired/dropped. */
public T? get<T>(string type_name);
}
Internally holds a type_name → slot_key map (captured at watch time from
request.held). On each get<T>, reads the CURRENT state from
state_service.get_slot(key). If the slot expired, returns null.
derive returns |
Action |
|---|---|
New TState instance |
state_service.update(slot_key, state) → SSE push |
current (same reference) |
No push (unchanged / not relevant) |
null |
state_service.clear_slot(slot_key) → SSE clear + remove watcher |
| Throws exception | Same as null: clear slot + remove watcher |
// Startup
statum.topic<CatTopic>("cat");
statum.topic_modifier<CatBreedDetailModifier, Cat, CatBreedDetail>("cat");
statum.topic_modifier<CatSummaryModifier, Cat, CatSummary>("cat");
statum.topic<T>(prefix) registers the topic handler for a key prefix (e.g.
"cat" matches "cat:42", "cat:99").
statum.topic_modifier<M, TPayload, TState>(prefix) registers a state modifier
for the prefix. Multiple modifiers can share a prefix (one per slot type).
public class CatDetailEntrypoint : StatumEntrypoint {
private TopicRegistry topics = inject<TopicRegistry>();
public override async DirectiveBuilder handle() throws GLib.Error {
var cat = yield db.find_cat_by_breed(request.route("breed"));
var detail_slot = state_service.new_slot(Scope.PAGE, to_detail_state(cat));
var summary_slot = state_service.new_slot(Scope.PAGE, to_summary_state(cat));
// Capture the full held context (type_name → slot_key) so modifiers
// can query auth, related entities, etc. at trigger time.
topics.watch(@"cat:$(cat.id)", detail_slot.id, typeof(CatBreedDetail).name(), request.held);
topics.watch(@"cat:$(cat.id)", summary_slot.id, typeof(CatSummary).name(), request.held);
return directives()
.set(detail_slot)
.set(summary_slot)
.subscribe(detail_slot.id)
.subscribe(summary_slot.id);
}
}
// From any context (worker, webhook, admin action)
yield topics.trigger("cat:42");
CatTopic).topic_key. If none, skip the load
entirely — no database query, no DI construction. The trigger is a no-op.topic.load("cat:42") → Cat (one DB query).state_type_name matches the watcher's slot type.
b. Read current from state_service.get_slot(watcher.slot_key).
c. Build HeldStates from watcher.held_keys.
d. Call modifier.derive(payload, current, held).
e. If result is null OR the call throws → clear_slot + remove watcher.
f. If result == current (reference equality) → skip.
g. If result is a new instance → update(slot_key, state) → SSE push.| Event | Action |
|---|---|
derive returns null |
clear_slot + remove watcher |
derive throws |
Same as null: clear_slot + remove watcher |
| SSE disconnect (navigate away) | update push is no-op; periodic sweep removes stale watchers |
| Server restart | Watchers lost from memory; re-registered on next entrypoint load |
| Slot expired from cache | trigger catches SLOT_NOT_FOUND in update → removes watcher |
| Periodic sweep (every 5 min) | Checks get_slot(key) for each watcher; removes if null |
HeldStates — src/TopicSystem.vala (or src/HeldStates.vala)public class HeldStates : Object {
private Dictionary<string, string> slot_keys;
private StateService state_service;
internal HeldStates(Dictionary<string, string> slot_keys, StateService state_service) { ... }
public T? get<T>(string type_name) { ... }
public Properties? get_properties(string type_name) { ... }
}
get<T> reads the current state from state_service.get_slot, maps to T via
GObjectMapping.from_properties_typed<T>.
Topic<TPayload> base class — src/TopicSystem.valapublic abstract class Topic<TPayload> : Object {
public abstract async TPayload load(string topic_key) throws GLib.Error;
}
StateModifier<TPayload, TState> + non-generic bridge — src/TopicSystem.valaThe generic class stores the payload/state types. A non-generic internal
interface (TopicStateModifier) lets the registry store heterogeneous modifiers:
internal interface TopicStateModifier : Object {
public abstract string state_type_name { get; }
public abstract async ModifierOutcome apply(Object payload, string slot_key, HeldStates held) throws GLib.Error;
}
public enum ModifierOutcome { UPDATED, UNCHANGED, CLEARED }
public abstract class StateModifier<TPayload, TState> : Object, TopicStateModifier {
protected StateService state_service = inject<StateService>();
public string state_type_name { get { return typeof(TState).name(); } }
public abstract TState? derive(TPayload payload, TState? current, HeldStates held);
public async ModifierOutcome apply(Object payload, string slot_key, HeldStates held) throws GLib.Error {
var slot = state_service.get_slot(slot_key);
TState? current = null;
if (slot != null && slot.current_state != null) {
try {
current = GObjectMapping.from_properties_typed<TState>(typeof(TState), slot.current_state.public_data);
} catch { current = null; }
}
var result = derive((TPayload) payload, current, held);
if (result == null) {
yield state_service.clear_slot(slot_key);
return ModifierOutcome.CLEARED;
}
if (result == current) {
return ModifierOutcome.UNCHANGED;
}
var new_state = new State() {
type_name = state_type_name,
public_data = GObjectMapping.to_properties((Object) result),
private_data = slot?.current_state?.private_data ?? new PropertyDictionary()
};
yield state_service.update(slot_key, new_state);
return ModifierOutcome.UPDATED;
}
}
TopicRegistry (singleton) — src/TopicSystem.valapublic class TopicRegistry : Object {
private StateService state_service = inject<StateService>();
// topic_key → watchers
private Dictionary<string, Series<Watcher>> topic_watchers;
// topic prefix → Topic handler (stored as a factory, resolved via DI)
private Dictionary<string, Type> topic_types;
// topic prefix → modifiers (by state_type_name)
private Dictionary<string, Dictionary<string, TopicStateModifier>> prefix_modifiers;
// Watcher record
private class Watcher : Object {
public string slot_key;
public string slot_type;
public Dictionary<string, string> held_keys; // type_name → slot_key
}
public void watch(string topic_key, string slot_key, string slot_type,
ReadOnlyAssociative<string, HeldSlot> request_held) { ... }
public async void trigger(string topic_key) throws GLib.Error {
// Early exit: if nobody is watching, skip the load entirely (no DB query).
var watchers = topic_watchers.get_or_default(topic_key);
if (watchers == null || ((!)watchers).to_array().length == 0) return;
var prefix = extract_prefix(topic_key); // "cat:42" → "cat"
var topic = construct_topic(prefix);
var payload = yield topic.load(topic_key);
var to_remove = new Series<Watcher>();
foreach (var watcher in (!)watchers) {
var modifiers = prefix_modifiers.get_or_default(prefix);
if (modifiers == null) continue;
TopicStateModifier modifier;
if (!((!)modifiers).try_get(watcher.slot_type, out modifier)) continue;
var held = new HeldStates(watcher.held_keys, state_service);
try {
var outcome = yield modifier.apply(payload, watcher.slot_key, held);
if (outcome == ModifierOutcome.CLEARED) {
to_remove.add(watcher);
}
} catch (GLib.Error e) {
// Exception in derive → clear + remove (same as null return)
try { yield state_service.clear_slot(watcher.slot_key); }
catch {}
to_remove.add(watcher);
}
}
foreach (var w in to_remove) {
((!)watchers).remove(w);
}
}
// Periodic cleanup: remove watchers whose slots are no longer cached.
public void sweep() { ... }
}
StatumConfigurator — register topics + modifierspublic void topic<TTopic>(string prefix) throws Error {
container.register_transient<TTopic>();
// Store typeof(TTopic) for the prefix
}
public void topic_modifier<TModifier, TPayload, TState>(string prefix) throws Error {
container.register_transient<TModifier>();
// Store typeof(TModifier) for the prefix + typeof(TState).name()
}
The registry resolves modifiers and topics from the container (DI-constructed).
StatumModule — register TopicRegistry singleton + sweep timercontainer.register_singleton<TopicRegistry>();
// Sweep timer runs every 5 minutes (background).
Both StatumEntrypoint and StatumAction gain:
protected TopicRegistry topics = inject<TopicRegistry>();
Add a CatTopic + CatBreedDetailModifier to the example (or a simpler variant
using the existing announcement/counter models to demonstrate the flow without a
database).
model.md: add a "Topic system" section.getting-started.md: add Topic<T>, StateModifier<T,S>, TopicRegistry
to the Vala API reference.// Topic: loads the entity
public class CatTopic : Topic<Cat> {
private Database db = inject<Database>();
public override async Cat load(string topic_key) throws GLib.Error {
return yield db.load_cat(int.parse(topic_key.split(":")[1]));
}
}
// Modifier: derives CatBreedDetail from Cat + held context
public class CatBreedDetailModifier : StateModifier<Cat, CatBreedDetail> {
public override CatBreedDetail? derive(Cat cat, CatBreedDetail? current, HeldStates held) {
var auth = held.get<AuthPublic>("auth");
if (auth == null) return null; // expired → clear
if (current != null && current.breed != cat.breed)
return current; // different breed → skip
return new CatBreedDetail() { ... }; // push
}
}
// Registration
statum.topic<CatTopic>("cat");
statum.topic_modifier<CatBreedDetailModifier, Cat, CatBreedDetail>("cat");
// Entrypoint
topics.watch(@"cat:$(id)", slot.id, typeof(CatBreedDetail).name(), request.held);
// Trigger (anywhere)
yield topics.trigger("cat:42");
load() (DB) + per-watcher: one get_slot + one derive (in-memory) + optionally one update (sign + push).get_slot. O(N) where N = watcher count.