Skip to content

feat(streamable_http): per-subscriber metadata and targeted sends - #218

Merged
zoedsoupe merged 2 commits into
zoedsoupe:mainfrom
fypm-labs:pr/streamable-http-subscriber-metadata
Jul 16, 2026
Merged

zoedsoupe merged 2 commits into
zoedsoupe:mainfrom
fypm-labs:pr/streamable-http-subscriber-metadata

Conversation

@BobbieBarker

Copy link
Copy Markdown
Contributor

Problem

The Streamable HTTP transport can push a message to a single session's SSE stream (route_to_session/3) or broadcast to every connected handler (send_message/3), but there is no way to deliver to a selected subset of connected subscribers, and no way for a host to attach application context to a subscriber or observe connected-handler counts. This is the single-node building block behind the cross-node delivery discussed in #189 / #190.

Solution

Add opaque per-subscriber metadata plus selector-based delivery:

  • register_sse_handler/3 stores an opaque metadata map verbatim; the 2-arity form delegates with %{}. The transport never interprets the map. Handlers stay keyed by session_id, so get/route/unregister remain O(1); the stored value widens from {pid, ref} to {pid, ref, metadata}.
  • handler_count/1 (total) and handler_count/2 (a predicate over metadata) expose connected-handler gauges.
  • send_message_to_subscribers/4 fans a message out to every handler whose metadata satisfies a caller-supplied selector -- filling the gap between route_to_session/3 (one session) and send_message/3 (all handlers).
  • The Plug gains a :subscriber_metadata option -- a (Plug.Conn.t() -> map()) callback invoked when an SSE stream opens -- so hosts can tag subscribers from the request (tenant, user, feature scope). Defaults to %{}, so existing behavior is unchanged.

Rationale

Server-initiated delivery to an application-defined subset of subscribers is a general MCP-server need (see #189 / #190). The library cannot know a host's routing dimensions, so metadata is kept fully opaque and host-populated via the plug callback; the transport provides only storage + selector primitives. Keeping the session_id key (rather than re-keying by {session_id, pid}) preserves the existing O(1) lookups. This is not spec-mandated -- it is a capability extension layered on the existing per-session SSE registry -- and I'm happy to adjust the shape (e.g. the callback signature) to your preference.

Tests

  • Transport: register_sse_handler/3 metadata round-trip + handler_count/1,2; 2-arg default of %{}; send_message_to_subscribers/4 delivers only to matching subscribers.
  • Plug: init/1 stores and defaults the :subscriber_metadata callback; end-to-end GET SSE registration attaches the configured metadata (verified via handler_count/2).

Add the ability to attach opaque, host-defined metadata to each SSE
subscriber and to deliver a server-to-client message to an arbitrary
subset of connected subscribers -- filling the gap between
`route_to_session/3` (one session) and `send_message/3` (broadcast to all).

- `register_sse_handler/3` stores an opaque metadata map verbatim; the
  2-arity form delegates with `%{}`. The transport never interprets the
  map. Handlers stay keyed by `session_id`, so get/route/unregister remain
  O(1); the stored value widens from `{pid, ref}` to `{pid, ref, metadata}`.
- `handler_count/1` and `handler_count/2` (a predicate over metadata)
  expose connected-handler gauges.
- `send_message_to_subscribers/4` fans a message out to every handler whose
  metadata satisfies a caller-supplied selector.
- The Plug gains a `:subscriber_metadata` option -- a
  `(Plug.Conn.t() -> map())` callback invoked when an SSE stream opens --
  so hosts can tag subscribers from the request (tenant, user, feature
  scope). Defaults to `%{}`.

Hosts stay in control of what the metadata means; the library provides
only the storage and selector primitives.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
@coderabbitai

coderabbitai Bot commented Jul 15, 2026 •

Copy link
Copy Markdown

Warning

Review limit reached

@zoedsoupe, you've reached your PR review limit, so we couldn't start this review.

Next review available in: 40 minutes

Enable usage-based reviews in Billing to review now. Otherwise, wait until the next included review is available.
You're only billed for reviews past your plan's rate limits ($0.25/file).

How can I continue?

After more reviews become available, a review can be triggered using the @coderabbitai review command as a PR comment. Alternatively, push new commits to this PR.

To avoid repeated limits, reduce automatic review volume by pausing incremental auto-reviews earlier, using label-based review opt-in, excluding WIP or generated PR titles, or requesting reviews manually when the PR is ready. If your team needs uninterrupted high-volume reviews, an organization admin can enable usage-based reviews.

How do review limits work?

CodeRabbit enforces per-developer PR review limits for each organization. Most developers receive the normal plan review availability.

For paid Pro and Pro+ PR reviews, CodeRabbit uses adaptive limits for sustained high-volume activity. When a developer's recent PR review activity reaches the 95th percentile or higher among CodeRabbit users, additional reviews become available more gradually as earlier reviews age out of the rolling window.

Please refer docs for additional details.

Review details
⚙️ Run configuration

Configuration used: Repository UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: 2c725d67-1609-4539-afd0-53e7d08217fa

📥 Commits

Reviewing files that changed from the base of the PR and between b32134a and 221c1ff.

📒 Files selected for processing (4)
  • lib/anubis/server/transport/streamable_http.ex
  • lib/anubis/server/transport/streamable_http/plug.ex
  • test/anubis/server/transport/streamable_http/plug_test.exs
  • test/anubis/server/transport/streamable_http_test.exs
✨ Finishing Touches
✨ Simplify code
  • Create PR with simplified code

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@zoedsoupe

Copy link
Copy Markdown
Owner

@coderabbitai review

@coderabbitai

coderabbitai Bot commented Jul 16, 2026 •

Copy link
Copy Markdown
✅ Action performed

Review finished.

Note: CodeRabbit is an incremental review system and does not re-review already reviewed commits. This command is applicable only when automatic reviews are paused.

@zoedsoupe
zoedsoupe merged commit c606658 into zoedsoupe:main Jul 16, 2026
12 checks passed
@zoedsoupe zoedsoupe mentioned this pull request Jul 16, 2026
zoedsoupe added a commit that referenced this pull request Jul 16, 2026
🚀 Want to release this?
---


##
[1.7.0](v1.6.2...v1.7.0)
(2026-07-16)


### Features

* **streamable_http:** per-subscriber metadata and targeted sends
([#218](#218))
([c606658](c606658))


### Bug Fixes

* Forward configured :headers on the DELETE session-teardown request
(follow-up to
[#180](#180))
([#213](#213))
([b32134a](b32134a))
* prevent "Server not initialized" race on first request
([#198](#198))
([e84624c](e84624c))
* **server:** resolve session names via Registry to prevent
atom-exhaustion DoS
([#188](#188))
([17e4a6d](17e4a6d))
* **session:** trap_exit so terminate/2 runs on supervisor shutdown
([#209](#209))
([6335cf4](6335cf4))
* **streamable_http:** don't close superseded SSE handler to prevent
reconnect flap
([#215](#215))
([a1e0ce6](a1e0ce6))


### Continuous Integration

* add new elixir versions
([3b636a8](3b636a8))
* add pr-quality workflow
([c0ca08f](c0ca08f))
* fix zig correct version for burrito
([6e410bd](6e410bd))
* use mlugg/setup-zig 0.15.2 in release-please auto build job
([2ed6187](2ed6187))

---
This PR was generated with [Release
Please](https://github.com/googleapis/release-please). See
[documentation](https://github.com/googleapis/release-please#release-please).
zoedsoupe added a commit that referenced this pull request Jul 16, 2026
## Problem

The Streamable HTTP transport can push a message to a single session's
SSE stream (`route_to_session/3`) or broadcast to every connected
handler (`send_message/3`), but there is no way to deliver to a
*selected subset* of connected subscribers, and no way for a host to
attach application context to a subscriber or observe connected-handler
counts. This is the single-node building block behind the cross-node
delivery discussed in #189 / #190.

## Solution

Add opaque per-subscriber metadata plus selector-based delivery:

- `register_sse_handler/3` stores an opaque `metadata` map verbatim; the
2-arity form delegates with `%{}`. The transport never interprets the
map. Handlers stay keyed by `session_id`, so `get`/`route`/`unregister`
remain O(1); the stored value widens from `{pid, ref}` to `{pid, ref,
metadata}`.
- `handler_count/1` (total) and `handler_count/2` (a predicate over
metadata) expose connected-handler gauges.
- `send_message_to_subscribers/4` fans a message out to every handler
whose metadata satisfies a caller-supplied selector -- filling the gap
between `route_to_session/3` (one session) and `send_message/3` (all
handlers).
- The Plug gains a `:subscriber_metadata` option -- a `(Plug.Conn.t() ->
map())` callback invoked when an SSE stream opens -- so hosts can tag
subscribers from the request (tenant, user, feature scope). Defaults to
`%{}`, so existing behavior is unchanged.

## Rationale

Server-initiated delivery to an application-defined subset of
subscribers is a general MCP-server need (see #189 / #190). The library
cannot know a host's routing dimensions, so metadata is kept fully
opaque and host-populated via the plug callback; the transport provides
only storage + selector primitives. Keeping the `session_id` key (rather
than re-keying by `{session_id, pid}`) preserves the existing O(1)
lookups. This is not spec-mandated -- it is a capability extension
layered on the existing per-session SSE registry -- and I'm happy to
adjust the shape (e.g. the callback signature) to your preference.

## Tests

- Transport: `register_sse_handler/3` metadata round-trip +
`handler_count/1,2`; 2-arg default of `%{}`;
`send_message_to_subscribers/4` delivers only to matching subscribers.
- Plug: `init/1` stores and defaults the `:subscriber_metadata`
callback; end-to-end GET SSE registration attaches the configured
metadata (verified via `handler_count/2`).


Co-authored-by: zoey <zoey.spessanha@zeetech.io>
zoedsoupe added a commit that referenced this pull request Jul 16, 2026
🚀 Want to release this?
---


##
[1.7.0](v1.6.2...v1.7.0)
(2026-07-16)


### Features

* **streamable_http:** per-subscriber metadata and targeted sends
([#218](#218))
([d6cf7b1](d6cf7b1))


### Bug Fixes

* Forward configured :headers on the DELETE session-teardown request
(follow-up to
[#180](#180))
([#213](#213))
([4edd2c0](4edd2c0))
* prevent "Server not initialized" race on first request
([#198](#198))
([6bb60f9](6bb60f9))
* **server:** resolve session names via Registry to prevent
atom-exhaustion DoS
([#188](#188))
([fdbc238](fdbc238))
* **session:** trap_exit so terminate/2 runs on supervisor shutdown
([#209](#209))
([a224c00](a224c00))
* **streamable_http:** don't close superseded SSE handler to prevent
reconnect flap
([#215](#215))
([e1cc4a8](e1cc4a8))


### Continuous Integration

* add new elixir versions
([26dd267](26dd267))
* add pr-quality workflow
([d67eaa9](d67eaa9))
* fix zig correct version for burrito
([afec768](afec768))
* use mlugg/setup-zig 0.15.2 in release-please auto build job
([428cace](428cace))

---
This PR was generated with [Release
Please](https://github.com/googleapis/release-please). See
[documentation](https://github.com/googleapis/release-please#release-please).
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants