mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 19:53:08 +00:00
feat(napplet): broker the NAP inc topic bus
Adds the `inc` pub/sub bus so napplets that declare it boot and can exchange
topic events: `inc.subscribe {topic}` / `inc.unsubscribe {topic}` register
interest, `inc.emit {topic, payload}` fans out an `inc.event {topic, payload,
sender}` to OTHER subscribed napplet sessions (never echoing the sender) — the
kehto runtime's inc contract.
- Router edge ops (gated on the INC declaration alone, like identity.watch —
no per-call consent): SubscribeInc/UnsubscribeInc/EmitInc outcomes.
- Protocol: readTopic/readPayloadRaw + encodeIncEvent.
- NappletIncBus in the broker service routes across the live napplet sessions
(the one service every sandbox binds), keyed by reply Messenger.
- Tests for inc routing + declaration gating; updated capability/router tests
that asserted the old "inc/theme/notify are unknown" behavior.
NOTE: napplets run foreground-only/one-at-a-time, so cross-napplet delivery is
usually a no-op in practice; the bus is correct if sessions ever overlap. It is
app-wide (not author-scoped) — a future refinement could namespace topics by
author. Unblocks feed/profile-viewer/chat/bot. See the plan doc.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016ncMHuBBVHEf7spAoSssde
This commit is contained in:
@@ -81,6 +81,9 @@ class NappletBrokerService : Service() {
|
||||
// Live relay subscriptions, keyed by the applet's subId; reads the current account live.
|
||||
private val liveSubscriptions = NappletLiveSubscriptions { Amethyst.instance.sessionManager.loggedInAccount() }
|
||||
|
||||
// The app-wide inc pub/sub bus: routes inc.emit between live napplet sessions as inc.event pushes.
|
||||
private val incBus = NappletIncBus { replyTo, payload -> push(replyTo, payload) }
|
||||
|
||||
// Streams identity.changed pushes (account switch / connect / disconnect) to a watching applet.
|
||||
private val identityWatch =
|
||||
NappletIdentityWatch(scope) {
|
||||
@@ -149,6 +152,9 @@ class NappletBrokerService : Service() {
|
||||
is NappletRequestRouter.Outcome.WatchIdentity -> identityWatch.start { push(replyTo, it) }
|
||||
is NappletRequestRouter.Outcome.UnwatchIdentity -> identityWatch.stop()
|
||||
is NappletRequestRouter.Outcome.Push -> outcome.payloads.forEach { push(replyTo, it) }
|
||||
is NappletRequestRouter.Outcome.SubscribeInc -> incBus.subscribe(replyTo, outcome.topic)
|
||||
is NappletRequestRouter.Outcome.UnsubscribeInc -> incBus.unsubscribe(replyTo, outcome.topic)
|
||||
is NappletRequestRouter.Outcome.EmitInc -> incBus.emit(replyTo, identity.coordinate, outcome.topic, outcome.payloadRaw)
|
||||
}
|
||||
}
|
||||
return true
|
||||
|
||||
@@ -0,0 +1,89 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.amethyst.napplet
|
||||
|
||||
import android.os.Messenger
|
||||
import com.vitorpamplona.amethyst.commons.napplet.protocol.NappletProtocolJson
|
||||
|
||||
/**
|
||||
* The NAP `inc` bus: a topic pub/sub that fans an `inc.emit` out to **other** napplet sessions
|
||||
* subscribed to that topic as an `inc.event` push (never echoing back to the sender). Lives in the
|
||||
* main-process broker service, which is the single point every sandbox session binds to, so it can
|
||||
* route between them.
|
||||
*
|
||||
* Subscribers are keyed by their reply [Messenger] (one per live napplet host). A napplet must have
|
||||
* declared the `inc` capability to subscribe or emit (enforced by the router before we get here).
|
||||
*
|
||||
* NOTE: Amethyst runs napplets **foreground-only, one at a time**, so in practice two sessions rarely
|
||||
* overlap and cross-napplet delivery is usually a no-op — but the bus is correct if they ever do.
|
||||
* The bus is app-wide (not scoped per author), so it deliberately allows different napplets to talk;
|
||||
* a future refinement could namespace topics by author if cross-napplet isolation is wanted.
|
||||
*/
|
||||
class NappletIncBus(
|
||||
private val deliver: (Messenger, String) -> Unit,
|
||||
) {
|
||||
// topic -> the reply Messengers of the napplet sessions subscribed to it.
|
||||
private val subscribers = HashMap<String, MutableSet<Messenger>>()
|
||||
|
||||
@Synchronized
|
||||
fun subscribe(
|
||||
who: Messenger,
|
||||
topic: String,
|
||||
) {
|
||||
subscribers.getOrPut(topic) { mutableSetOf() }.add(who)
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
fun unsubscribe(
|
||||
who: Messenger,
|
||||
topic: String,
|
||||
) {
|
||||
subscribers[topic]?.let { set ->
|
||||
set.remove(who)
|
||||
if (set.isEmpty()) subscribers.remove(topic)
|
||||
}
|
||||
}
|
||||
|
||||
/** Delivers an `inc.event` for [topic] to every subscriber except [from] (no self-echo). */
|
||||
@Synchronized
|
||||
fun emit(
|
||||
from: Messenger,
|
||||
sender: String,
|
||||
topic: String,
|
||||
payloadRaw: String,
|
||||
) {
|
||||
val set = subscribers[topic] ?: return
|
||||
if (set.isEmpty()) return
|
||||
val envelope = NappletProtocolJson.encodeIncEvent(topic, payloadRaw, sender)
|
||||
set.filter { it != from }.forEach { deliver(it, envelope) }
|
||||
}
|
||||
|
||||
/** Drops [who] from every topic (e.g. when its napplet host goes away). */
|
||||
@Synchronized
|
||||
fun removeAll(who: Messenger) {
|
||||
val empties = mutableListOf<String>()
|
||||
for ((topic, set) in subscribers) {
|
||||
set.remove(who)
|
||||
if (set.isEmpty()) empties.add(topic)
|
||||
}
|
||||
empties.forEach { subscribers.remove(it) }
|
||||
}
|
||||
}
|
||||
+3
-1
@@ -37,12 +37,14 @@ class NappletCapabilityTest {
|
||||
assertEquals(NappletCapability.STORAGE, NappletCapability.fromNapDomain(" STORAGE "))
|
||||
assertEquals(NappletCapability.RESOURCE, NappletCapability.fromNapDomain("resource"))
|
||||
assertEquals(NappletCapability.UPLOAD, NappletCapability.fromNapDomain("upload"))
|
||||
assertEquals(NappletCapability.THEME, NappletCapability.fromNapDomain("theme"))
|
||||
assertEquals(NappletCapability.NOTIFY, NappletCapability.fromNapDomain("notify"))
|
||||
assertEquals(NappletCapability.INC, NappletCapability.fromNapDomain("inc"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun unknownDomainMapsToNullNotAFallbackGrant() {
|
||||
// Domains we don't broker yet must stay unknown (default-deny), not fall through.
|
||||
assertNull(NappletCapability.fromNapDomain("inc"))
|
||||
assertNull(NappletCapability.fromNapDomain("intent"))
|
||||
assertNull(NappletCapability.fromNapDomain("cvm"))
|
||||
assertNull(NappletCapability.fromNapDomain("filesystem"))
|
||||
|
||||
+34
@@ -65,6 +65,22 @@ object NappletRequestRouter {
|
||||
data class Push(
|
||||
val payloads: List<String>,
|
||||
) : Outcome
|
||||
|
||||
/** Subscribe this napplet to the inc bus [topic]; it then receives `inc.event` pushes (fire-and-forget). */
|
||||
data class SubscribeInc(
|
||||
val topic: String,
|
||||
) : Outcome
|
||||
|
||||
/** Stop receiving inc-bus [topic] events (fire-and-forget). */
|
||||
data class UnsubscribeInc(
|
||||
val topic: String,
|
||||
) : Outcome
|
||||
|
||||
/** Emit [payloadRaw] on inc-bus [topic]; the host fans it out as `inc.event` to other subscribers. */
|
||||
data class EmitInc(
|
||||
val topic: String,
|
||||
val payloadRaw: String,
|
||||
) : Outcome
|
||||
}
|
||||
|
||||
suspend fun route(
|
||||
@@ -89,6 +105,24 @@ object NappletRequestRouter {
|
||||
return if (NappletCapability.IDENTITY in declared) Outcome.WatchIdentity else Outcome.Ignore
|
||||
"identity.unwatch" ->
|
||||
return Outcome.UnwatchIdentity
|
||||
// inc bus: a topic pub/sub between napplets/services, authorized on the INC declaration
|
||||
// alone (like identity.watch) — no per-call consent. The host owns the cross-session fan-out.
|
||||
"inc.subscribe" -> {
|
||||
val topic = runCatching { NappletProtocolJson.readTopic(payload) }.getOrNull()
|
||||
return if (topic != null && NappletCapability.INC in declared) Outcome.SubscribeInc(topic) else Outcome.Ignore
|
||||
}
|
||||
"inc.unsubscribe" -> {
|
||||
val topic = runCatching { NappletProtocolJson.readTopic(payload) }.getOrNull()
|
||||
return if (topic != null) Outcome.UnsubscribeInc(topic) else Outcome.Ignore
|
||||
}
|
||||
"inc.emit" -> {
|
||||
val topic = runCatching { NappletProtocolJson.readTopic(payload) }.getOrNull()
|
||||
return if (topic != null && NappletCapability.INC in declared) {
|
||||
Outcome.EmitInc(topic, runCatching { NappletProtocolJson.readPayloadRaw(payload) }.getOrDefault("null"))
|
||||
} else {
|
||||
Outcome.Ignore
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
val request =
|
||||
|
||||
+22
@@ -80,6 +80,28 @@ object NappletProtocolJson {
|
||||
put("subId", subId)
|
||||
}.toString()
|
||||
|
||||
/** The `topic` of an `inc.*` envelope (the pub/sub topic the napplet emits to or subscribes on). */
|
||||
fun readTopic(envelopeJson: String): String? = json.parseToJsonElement(envelopeJson).jsonObject.str("topic")
|
||||
|
||||
/** The raw `payload` JSON element of an `inc.emit`, serialized back to a string (`"null"` if absent). */
|
||||
fun readPayloadRaw(envelopeJson: String): String = (json.parseToJsonElement(envelopeJson).jsonObject["payload"] ?: JsonNull).toString()
|
||||
|
||||
/**
|
||||
* An `inc.event` push: delivers [topic]'s [payloadRaw] (a raw JSON value) to a subscriber, tagged
|
||||
* with the emitting napplet's [sender] coordinate. Mirrors the kehto runtime's inc-bus contract.
|
||||
*/
|
||||
fun encodeIncEvent(
|
||||
topic: String,
|
||||
payloadRaw: String,
|
||||
sender: String,
|
||||
): String =
|
||||
buildJsonObject {
|
||||
put("type", "inc.event")
|
||||
put("topic", topic)
|
||||
put("payload", runCatching { json.parseToJsonElement(payloadRaw) }.getOrDefault(JsonNull))
|
||||
put("sender", sender)
|
||||
}.toString()
|
||||
|
||||
/** A `keys.action` push: the user triggered the registered keyboard/command action [actionId]. */
|
||||
fun encodeKeysAction(actionId: String): String =
|
||||
buildJsonObject {
|
||||
|
||||
+25
-1
@@ -92,10 +92,34 @@ class NappletRequestRouterTest {
|
||||
assertTrue(outcome.payload.contains("resource.cancel.result"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun incSubscribeAndEmitBecomeBusOps() =
|
||||
runTest {
|
||||
assertEquals(
|
||||
NappletRequestRouter.Outcome.SubscribeInc("news"),
|
||||
route("""{"type":"inc.subscribe","topic":"news"}"""),
|
||||
)
|
||||
assertEquals(
|
||||
NappletRequestRouter.Outcome.UnsubscribeInc("news"),
|
||||
route("""{"type":"inc.unsubscribe","topic":"news"}"""),
|
||||
)
|
||||
val emit = route("""{"type":"inc.emit","topic":"news","payload":{"text":"hi"}}""")
|
||||
assertIs<NappletRequestRouter.Outcome.EmitInc>(emit)
|
||||
assertEquals("news", emit.topic)
|
||||
assertTrue(emit.payloadRaw.contains("hi"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun incIsIgnoredWhenNotDeclared() =
|
||||
runTest {
|
||||
val outcome = NappletRequestRouter.route(broker(), applet, emptySet(), """{"type":"inc.emit","topic":"t","payload":1}""")
|
||||
assertEquals(NappletRequestRouter.Outcome.Ignore, outcome)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun malformedRequestRepliesFailed() =
|
||||
runTest {
|
||||
val outcome = route("""{"type":"inc.emit","topic":"t"}""")
|
||||
val outcome = route("""{"type":"totally.unknown"}""")
|
||||
assertIs<NappletRequestRouter.Outcome.Reply>(outcome)
|
||||
assertTrue(outcome.payload.contains("failed"))
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user