topic-system-plan.md 14 KB

Topic System Plan — Typed topics, state modifiers, and live held-state lookup

Overview

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).


Design

Topic<TPayload> — loads the entity once per trigger

public 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 watcher

public 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 time

public 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.

Return-value semantics

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

Registration

// 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).

Entrypoint — watch with held context

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);
    }
}

Trigger

// From any context (worker, webhook, admin action)
yield topics.trigger("cat:42");

What the registry does on trigger

  1. Resolve the topic handler by prefix ("cat" → CatTopic).
  2. Check if there are any watchers for topic_key. If none, skip the load entirely — no database query, no DI construction. The trigger is a no-op.
  3. topic.load("cat:42")Cat (one DB query).
  4. Find all modifiers for the prefix.
  5. For each watcher: a. Find the modifier whose 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.

Cleanup

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

Tasks

V-1: HeldStatessrc/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>.

V-2: Topic<TPayload> base class — src/TopicSystem.vala

public abstract class Topic<TPayload> : Object {
    public abstract async TPayload load(string topic_key) throws GLib.Error;
}

V-3: StateModifier<TPayload, TState> + non-generic bridge — src/TopicSystem.vala

The 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;
    }
}

V-4: TopicRegistry (singleton) — src/TopicSystem.vala

public 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() { ... }
}

V-5: StatumConfigurator — register topics + modifiers

public 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).

V-6: StatumModule — register TopicRegistry singleton + sweep timer

container.register_singleton<TopicRegistry>();
// Sweep timer runs every 5 minutes (background).

V-7: Entrypoint/action access

Both StatumEntrypoint and StatumAction gain:

protected TopicRegistry topics = inject<TopicRegistry>();

V-8: Example

Add a CatTopic + CatBreedDetailModifier to the example (or a simpler variant using the existing announcement/counter models to demonstrate the flow without a database).

D-1: Documentation

  • model.md: add a "Topic system" section.
  • getting-started.md: add Topic<T>, StateModifier<T,S>, TopicRegistry to the Vala API reference.

API sketch (end-to-end)

// 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");

Server resources

  • Each watcher: ~150 bytes (topic key + slot key + type name + held keys map).
  • 10K concurrent watchers ≈ 1.5 MB.
  • Trigger cost: one load() (DB) + per-watcher: one get_slot + one derive (in-memory) + optionally one update (sign + push).
  • Sweep: every 5 minutes, iterates watchers, checks get_slot. O(N) where N = watcher count.

Out of scope

  • Per-watcher personalised payloads (different state for different watchers of the same topic). Use per-user topic keys instead.
  • Cross-process pub/sub (Redis/NATS backend for multi-server). Each process has its own TopicRegistry.
  • Pattern-based topic matching (wildcards). Prefix matching ("cat" matches "cat:42") is sufficient.