Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,8 @@ List<MentionCandidate> buildMentionCandidates({
avatarUrl: profile?.avatarUrl,
isAgent: true,
isMember: false,
ownerPubkey: ownerByAgentPubkey[pk] ?? profile?.ownerPubkey,
ownerPubkey:
agent.ownerPubkey ?? ownerByAgentPubkey[pk] ?? profile?.ownerPubkey,
),
);
}
Expand Down
152 changes: 152 additions & 0 deletions mobile/lib/shared/mentions/agent_discovery.dart
Original file line number Diff line number Diff line change
@@ -0,0 +1,152 @@
part of 'agent_identity_provider.dart';

// Coordinates seed discovery only. Exact signed profiles and current policy
// are resolved by readAgentAuthorization before any candidate is exposed.
Future<Set<String>> _ownedAgentKeys(
RelaySessionNotifier session,
String viewer,
bool Function() current,
) async {
final keys = <String>{};
NostrEvent? cursor;
for (var page = 0; page < 200; page++) {
if (!current()) throw StateError('Agent discovery scope changed');
final events = await session.queryRelay([
NostrFilter(
kinds: const [30177],
authors: [viewer],
limit: 500,
until: cursor?.createdAt,
extensions: {if (cursor != null) 'before_id': cursor.id},
),
]);
if (!current()) throw StateError('Agent discovery scope changed');
for (final event in events) {
final key = event.getTagValue('d');
if (event.kind == 30177 &&
event.pubkey == viewer &&
verifySignedEvent(event) &&
key != null &&
RegExp(r'^[0-9a-f]{64}$').hasMatch(key) &&
event.tags.where((tag) => tag.isNotEmpty && tag[0] == 'd').length ==
1) {
keys.add(key);
}
}
if (keys.length > 1000) {
throw StateError('Agent discovery exceeds key budget');
}
if (events.length < 500) return keys;
final next = events.last;
if (cursor != null &&
(next.createdAt > cursor.createdAt ||
(next.createdAt == cursor.createdAt &&
next.id.compareTo(cursor.id) <= 0))) {
throw StateError('Agent discovery pagination did not advance');
}
cursor = next;
}
throw StateError('Agent discovery exceeds page budget');
}

// Separate subscription owner: a refresh must not tear down/replay its trigger.
class _AgentDirectoryUpdates extends Notifier<int> {
Object? failure;
@override
int build() {
final sessionState = ref.watch(relaySessionProvider);
ref.watch(myPubkeyProvider);
ref.watch(relayConfigProvider);
failure = null;
var disposed = false;
var attempts = 0;
Timer? retry;
void Function()? unsubscribe;
Timer? debounce;
ref.onDispose(() {
disposed = true;
unsubscribe?.call();
debounce?.cancel();
retry?.cancel();
});
void changed() {
if (disposed || debounce != null) return;
debounce = Timer(const Duration(milliseconds: 150), () {
debounce = null;
if (!disposed) state++;
});
}

if (sessionState.status == SessionStatus.connected) {
final session = ref.read(relaySessionProvider.notifier);
Future<void> subscribe() async {
if (disposed) return;
attempts++;
var active = true;
bool current() => !disposed && active;
void lost(Object error) {
if (!current()) return;
active = false;
unsubscribe?.call();
unsubscribe = null;
failure = error;
changed();
// One budget for establishment and terminal closure per generation.
// Exhaustion remains visible until session/account rebuild.
if (attempts < 3) {
retry = Timer(Duration(milliseconds: 250 * attempts), () {
retry = null;
unawaited(subscribe());
});
}
}

try {
final close = await session.subscribeWithStatus(
NostrFilter(
kinds: const [0, 5, 10100, 30177, 39002],
limit: 0,
since: DateTime.now().millisecondsSinceEpoch ~/ 1000,
),
(event) {
if (!current()) return;
if (event.kind == 5 && !_isAgentCoordinateDeletion(event)) return;
changed();
},
onClosed: (message) => lost(StateError(message)),
onStatusChanged: (status) {
if (!current()) return;
if (status == RelaySubscriptionStatus.ready) failure = null;
changed();
},
);
if (!current()) {
close();
} else {
unsubscribe = close;
}
} catch (error) {
lost(error);
}
}

unawaited(subscribe());
}
return 0;
}
}

final _agentDirectoryUpdatesProvider =
NotifierProvider<_AgentDirectoryUpdates, int>(_AgentDirectoryUpdates.new);

// Desktop build_agent_delete emits a kind:5 a-coordinate, not a 30177 event.
bool _isAgentCoordinateDeletion(NostrEvent event) =>
verifySignedEvent(event) &&
event.tags.any((tag) {
if (tag.length < 2 || tag[0] != 'a') return false;
final coordinate = tag[1].split(':');
return coordinate.length == 3 &&
coordinate[0] == '30177' &&
coordinate[1] == event.pubkey &&
RegExp(r'^[0-9a-f]{64}$').hasMatch(coordinate[2]);
});
23 changes: 16 additions & 7 deletions mobile/lib/shared/mentions/agent_identity_provider.dart
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import 'dart:async';
import 'dart:collection';
import 'dart:convert';

Expand All @@ -10,6 +11,7 @@ import '../../shared/relay/relay.dart';

part 'agent_policy.dart';
part 'agent_authorization.dart';
part 'agent_discovery.dart';

/// A relay agent parsed from its kind:10100 agent-profile event.
///
Expand Down Expand Up @@ -65,27 +67,34 @@ Map<String, dynamic>? _tryDecodeJsonMap(String content) {
}
}

/// Relay agent directory from kind:10100 agent-profile events.
/// Relay directory, including verified owned agents without runtime records.
///
/// Watches the session and only fetches after the WebSocket connects.
final agentDirectoryProvider = FutureProvider<List<AgentDirectoryEntry>>((
ref,
) async {
ref.watch(_agentDirectoryUpdatesProvider);
final failure = ref.read(_agentDirectoryUpdatesProvider.notifier).failure;
if (failure != null) throw failure;
final sessionState = ref.watch(relaySessionProvider);
if (sessionState.status != SessionStatus.connected) return const [];
final session = ref.read(relaySessionProvider.notifier);
final config = ref.read(relayConfigProvider);
final config = ref.watch(relayConfigProvider);
var disposed = false;
ref.onDispose(() => disposed = true);
final viewer = ref.read(myPubkeyProvider);
final viewer = ref.watch(myPubkeyProvider);
bool current() =>
!disposed && identical(config, ref.read(relayConfigProvider));
if (viewer == null) return [];
final owned = await _ownedAgentKeys(session, viewer, current);
if (!current()) return [];
final events = await session.fetchHistory(NostrFilters.agentProfiles());
if (disposed) return [];
if (!current()) return [];
return readAgentAuthorization(
session,
events.map((event) => event.pubkey).toSet(),
{...owned, ...events.map((event) => event.pubkey)},
viewer: viewer,
isCurrent: () =>
!disposed && identical(config, ref.read(relayConfigProvider)),
isCurrent: current,
);
});

Expand Down
15 changes: 14 additions & 1 deletion mobile/test/features/channels/compose_bar_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,12 @@ import 'package:buzz/shared/widgets/anchored_popover_menu.dart';
import 'package:buzz/shared/widgets/mobile_tab_footer_backdrop.dart';
import 'package:shared_preferences/shared_preferences.dart';

import '../../shared/mentions/agent_policy_test.dart'
show PolicySession, signed;
import '../../shared/crypto/nip_oa_test.dart' show authTag, profile;

part 'discovery_lifecycle_tests.dart';

final _pngBytes = Uint8List.fromList([
0x89,
0x50,
Expand Down Expand Up @@ -174,6 +180,7 @@ Widget _buildComposeBar({
required ComposeBarOnSend onSend,
List<ChannelMember> members = const <ChannelMember>[],
Future<List<ChannelMember>>? membersFuture,
RelaySessionNotifier? discoverySession,
List<AgentDirectoryEntry> relayAgents = const <AgentDirectoryEntry>[],
List<Channel> channels = const <Channel>[],
List<ChannelMember> cachedMembers = const <ChannelMember>[],
Expand Down Expand Up @@ -210,7 +217,12 @@ Widget _buildComposeBar({
channelMembersProvider(
'channel-1',
).overrideWith((ref) => membersFuture ?? Future.value(members)),
agentDirectoryProvider.overrideWith((ref) async => relayAgents),
if (discoverySession == null)
agentDirectoryProvider.overrideWith((ref) async => relayAgents),
if (discoverySession != null)
relaySessionProvider.overrideWith(() => discoverySession),
if (discoverySession != null)
myPubkeyProvider.overrideWithValue(currentPubkey),
agentOwnersProvider.overrideWith((ref) async => const <String, String>{}),
relayClientProvider.overrideWithValue(
RelayClient(baseUrl: 'http://localhost:3000'),
Expand Down Expand Up @@ -668,6 +680,7 @@ void main() {
_setMockMediaUploadPlatformHandler(null);
});

discoveryLifecycleTests();
group('ComposeBar', () {
testWidgets('starts compact and grows to the full-width composer', (
tester,
Expand Down
89 changes: 89 additions & 0 deletions mobile/test/features/channels/discovery_lifecycle_tests.dart
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
part of 'compose_bar_test.dart';

void discoveryLifecycleTests() {
testWidgets(
'terminal CLOSED recovers; signed roster removal drops open picker',
(tester) async {
final viewer = nostr.Keys.generate();
final owner = nostr.Keys.generate();
final agent = nostr.Keys.generate();
final relay = nostr.Keys.generate();
NostrEvent roster(int time, bool includeViewer) => signed(
relay,
39002,
'',
time: time,
tags: [
['d', 'channel-1'],
if (includeViewer) ['p', viewer.public],
['p', agent.public, '', 'bot'],
],
);
final events = [
roster(100, true),
profile(agent, [authTag(owner, agent.public)]),
signed(agent, 10100, {
'name': 'Helper Bot',
'channel_ids': ['channel-1'],
}),
signed(
owner,
30177,
{'name': 'Helper Bot', 'parallelism': 1, 'respond_to': 'anyone'},
tags: [
['d', agent.public],
],
),
];
// PolicySession replaces queries only: live REQ/EOSE/CLOSED/event handling
// and subscription removal/disposal are the production relay SDK.
final session = PolicySession(events);
await tester.pumpWidget(
_buildComposeBar(
discoverySession: session,
currentPubkey: viewer.public,
uploadService: _testUploadService(viewer.nsec),
channels: [_makeCurrentChannel()],
onSend:
(
content,
mentions, {
mediaTags = const <List<String>>[],
}) async {},
),
);
await _expandComposer(tester);
await tester.enterText(find.byType(TextField), '@hel');
await tester.pump();
session.debugHandleMessage(['EOSE', 'l-1']);
await tester.pump(const Duration(milliseconds: 150));
await tester.pumpAndSettle();
final c = ProviderScope.containerOf(
tester.element(find.byType(ComposeBar)),
);
expect(c.read(agentDirectoryProvider).requireValue, hasLength(1));
expect(find.text('Helper Bot'), findsOneWidget);
session.debugHandleMessage(['CLOSED', 'l-1', 'restricted: terminal']);
await tester.pump(const Duration(milliseconds: 150));
await tester.pump();
expect(c.read(agentDirectoryProvider).hasError, isTrue);
expect(find.text('Helper Bot'), findsNothing);
await tester.pump(const Duration(milliseconds: 100));
session.debugHandleMessage(['EOSE', 'l-2']);
await tester.pump(const Duration(milliseconds: 150));
await tester.pumpAndSettle();
expect(find.text('Helper Bot'), findsOneWidget);
final removed = roster(101, false);
events.add(removed);
session.debugHandleMessage(['EVENT', 'l-2', removed.toJson()]);
session.debugFlushEventBuffer();
await tester.pump(const Duration(milliseconds: 150));
await tester.pumpAndSettle();
expect(c.read(agentDirectoryProvider).requireValue, isEmpty);
expect(find.text('Helper Bot'), findsNothing);
expect(find.byType(TextField), findsOneWidget); // same mounted editor
await tester.pumpWidget(const SizedBox.shrink());
session.debugDispose();
},
);
}
Loading