Make RelayPool private to NostrNetworkManager and migrate usages
Signed-off-by: Daniel D’Aquino <daniel@daquino.me>
This commit is contained in:
@@ -39,63 +39,41 @@ class SearchHomeModel: ObservableObject {
|
||||
self.objectWillChange.send()
|
||||
}
|
||||
|
||||
func subscribe() {
|
||||
func load() async {
|
||||
loading = true
|
||||
let to_relays = determine_to_relays(pool: damus_state.nostrNetwork.pool, filters: damus_state.relay_filters)
|
||||
|
||||
var follow_list_filter = NostrFilter(kinds: [.follow_list])
|
||||
follow_list_filter.until = UInt32(Date.now.timeIntervalSince1970)
|
||||
let to_relays = damus_state.nostrNetwork.ourRelayDescriptors
|
||||
.map { $0.url }
|
||||
.filter { !damus_state.relay_filters.is_filtered(timeline: .search, relay_id: $0) }
|
||||
|
||||
damus_state.nostrNetwork.pool.subscribe(sub_id: base_subid, filters: [get_base_filter()], handler: handle_event, to: to_relays)
|
||||
damus_state.nostrNetwork.pool.subscribe(sub_id: follow_pack_subid, filters: [follow_list_filter], handler: handle_event, to: to_relays)
|
||||
}
|
||||
|
||||
func unsubscribe(to: RelayURL? = nil) {
|
||||
loading = false
|
||||
damus_state.nostrNetwork.pool.unsubscribe(sub_id: base_subid, to: to.map { [$0] })
|
||||
damus_state.nostrNetwork.pool.unsubscribe(sub_id: follow_pack_subid, to: to.map { [$0] })
|
||||
}
|
||||
|
||||
func handle_event(relay_id: RelayURL, conn_ev: NostrConnectionEvent) {
|
||||
guard case .nostr_event(let event) = conn_ev else {
|
||||
return
|
||||
for await item in damus_state.nostrNetwork.reader.subscribe(filters: [get_base_filter()], to: to_relays) {
|
||||
switch item {
|
||||
case .event(let borrow):
|
||||
var event: NostrEvent? = nil
|
||||
try? borrow { ev in
|
||||
event = ev.toOwned()
|
||||
}
|
||||
guard let event else { return }
|
||||
await self.handleEvent(event)
|
||||
case .eose: break
|
||||
}
|
||||
}
|
||||
loading = false
|
||||
|
||||
switch event {
|
||||
case .event(let sub_id, let ev):
|
||||
guard sub_id == self.base_subid || sub_id == self.profiles_subid || sub_id == self.follow_pack_subid else {
|
||||
guard let txn = NdbTxn(ndb: damus_state.ndb) else { return }
|
||||
load_profiles(context: "universe", load: .from_events(events.all_events), damus_state: damus_state, txn: txn)
|
||||
}
|
||||
|
||||
@MainActor
|
||||
func handleEvent(_ ev: NostrEvent) {
|
||||
if ev.is_textlike && should_show_event(state: damus_state, ev: ev) && !ev.is_reply() {
|
||||
if !damus_state.settings.multiple_events_per_pubkey && seen_pubkey.contains(ev.pubkey) {
|
||||
return
|
||||
}
|
||||
if ev.is_textlike && should_show_event(state: damus_state, ev: ev) && !ev.is_reply()
|
||||
{
|
||||
if !damus_state.settings.multiple_events_per_pubkey && seen_pubkey.contains(ev.pubkey) {
|
||||
return
|
||||
}
|
||||
seen_pubkey.insert(ev.pubkey)
|
||||
|
||||
if self.events.insert(ev) {
|
||||
self.objectWillChange.send()
|
||||
}
|
||||
}
|
||||
case .notice(let msg):
|
||||
print("search home notice: \(msg)")
|
||||
case .ok:
|
||||
break
|
||||
case .eose(let sub_id):
|
||||
loading = false
|
||||
seen_pubkey.insert(ev.pubkey)
|
||||
|
||||
if sub_id == self.base_subid {
|
||||
// Make sure we unsubscribe after we've fetched the global events
|
||||
// global events are not realtime
|
||||
unsubscribe(to: relay_id)
|
||||
|
||||
guard let txn = NdbTxn(ndb: damus_state.ndb) else { return }
|
||||
load_profiles(context: "universe", profiles_subid: profiles_subid, relay_id: relay_id, load: .from_events(events.all_events), damus_state: damus_state, txn: txn)
|
||||
if self.events.insert(ev) {
|
||||
self.objectWillChange.send()
|
||||
}
|
||||
|
||||
break
|
||||
case .auth:
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -135,44 +113,35 @@ enum PubkeysToLoad {
|
||||
case from_keys([Pubkey])
|
||||
}
|
||||
|
||||
func load_profiles<Y>(context: String, profiles_subid: String, relay_id: RelayURL, load: PubkeysToLoad, damus_state: DamusState, txn: NdbTxn<Y>) {
|
||||
func load_profiles<Y>(context: String, load: PubkeysToLoad, damus_state: DamusState, txn: NdbTxn<Y>) {
|
||||
let authors = find_profiles_to_fetch(profiles: damus_state.profiles, load: load, cache: damus_state.events, txn: txn)
|
||||
|
||||
guard !authors.isEmpty else {
|
||||
return
|
||||
}
|
||||
|
||||
print("load_profiles[\(context)]: requesting \(authors.count) profiles from \(relay_id)")
|
||||
|
||||
let filter = NostrFilter(kinds: [.metadata], authors: authors)
|
||||
|
||||
damus_state.nostrNetwork.pool.subscribe_to(sub_id: profiles_subid, filters: [filter], to: [relay_id]) { rid, conn_ev in
|
||||
Task {
|
||||
print("load_profiles[\(context)]: requesting \(authors.count) profiles from relay pool")
|
||||
let filter = NostrFilter(kinds: [.metadata], authors: authors)
|
||||
|
||||
let now = UInt64(Date.now.timeIntervalSince1970)
|
||||
switch conn_ev {
|
||||
case .ws_connection_event:
|
||||
break
|
||||
case .nostr_event(let ev):
|
||||
guard ev.subid == profiles_subid, rid == relay_id else { return }
|
||||
|
||||
switch ev {
|
||||
case .event(_, let ev):
|
||||
if ev.known_kind == .metadata {
|
||||
damus_state.ndb.write_profile_last_fetched(pubkey: ev.pubkey, fetched_at: now)
|
||||
for await item in damus_state.nostrNetwork.reader.subscribe(filters: [filter]) {
|
||||
let now = UInt64(Date.now.timeIntervalSince1970)
|
||||
switch item {
|
||||
case .event(let borrow):
|
||||
var event: NostrEvent? = nil
|
||||
try? borrow { ev in
|
||||
event = ev.toOwned()
|
||||
}
|
||||
guard let event else { return }
|
||||
if event.known_kind == .metadata {
|
||||
damus_state.ndb.write_profile_last_fetched(pubkey: event.pubkey, fetched_at: now)
|
||||
}
|
||||
case .eose:
|
||||
print("load_profiles[\(context)]: done loading \(authors.count) profiles from \(relay_id)")
|
||||
damus_state.nostrNetwork.pool.unsubscribe(sub_id: profiles_subid, to: [relay_id])
|
||||
case .ok:
|
||||
break
|
||||
case .notice:
|
||||
break
|
||||
case .auth:
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
print("load_profiles[\(context)]: done loading \(authors.count) profiles from relay pool")
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -14,8 +14,8 @@ class SearchModel: ObservableObject {
|
||||
@Published var loading: Bool = false
|
||||
|
||||
var search: NostrFilter
|
||||
let sub_id = UUID().description
|
||||
let profiles_subid = UUID().description
|
||||
var listener: Task<Void, Never>? = nil
|
||||
let limit: UInt32 = 500
|
||||
|
||||
init(state: DamusState, search: NostrFilter) {
|
||||
@@ -39,17 +39,32 @@ class SearchModel: ObservableObject {
|
||||
search.kinds = [.text, .like, .longform, .highlight, .follow_list]
|
||||
|
||||
//likes_filter.ids = ref_events.referenced_ids!
|
||||
|
||||
print("subscribing to search '\(search)' with sub_id \(sub_id)")
|
||||
state.nostrNetwork.pool.register_handler(sub_id: sub_id, handler: handle_event)
|
||||
loading = true
|
||||
state.nostrNetwork.pool.send(.subscribe(.init(filters: [search], sub_id: sub_id)))
|
||||
listener?.cancel()
|
||||
listener = Task {
|
||||
self.loading = true
|
||||
print("subscribing to search")
|
||||
for await item in await state.nostrNetwork.reader.subscribe(filters: [search]) {
|
||||
switch item {
|
||||
case .event(let borrow):
|
||||
try? borrow { ev in
|
||||
let event = ev.toOwned()
|
||||
if event.is_textlike && event.should_show_event {
|
||||
self.add_event(event)
|
||||
}
|
||||
}
|
||||
case .eose:
|
||||
break
|
||||
}
|
||||
guard let txn = NdbTxn(ndb: state.ndb) else { return }
|
||||
load_profiles(context: "search", load: .from_events(self.events.all_events), damus_state: state, txn: txn)
|
||||
}
|
||||
self.loading = false
|
||||
}
|
||||
}
|
||||
|
||||
func unsubscribe() {
|
||||
state.nostrNetwork.pool.unsubscribe(sub_id: sub_id)
|
||||
loading = false
|
||||
print("unsubscribing from search '\(search)' with sub_id \(sub_id)")
|
||||
listener?.cancel()
|
||||
listener = nil
|
||||
}
|
||||
|
||||
func add_event(_ ev: NostrEvent) {
|
||||
@@ -65,25 +80,6 @@ class SearchModel: ObservableObject {
|
||||
objectWillChange.send()
|
||||
}
|
||||
}
|
||||
|
||||
func handle_event(relay_id: RelayURL, ev: NostrConnectionEvent) {
|
||||
let (sub_id, done) = handle_subid_event(pool: state.nostrNetwork.pool, relay_id: relay_id, ev: ev) { sub_id, ev in
|
||||
if ev.is_textlike && ev.should_show_event {
|
||||
self.add_event(ev)
|
||||
}
|
||||
}
|
||||
|
||||
guard done else {
|
||||
return
|
||||
}
|
||||
|
||||
self.loading = false
|
||||
|
||||
if sub_id == self.sub_id {
|
||||
guard let txn = NdbTxn(ndb: state.ndb) else { return }
|
||||
load_profiles(context: "search", profiles_subid: self.profiles_subid, relay_id: relay_id, load: .from_events(self.events.all_events), damus_state: state, txn: txn)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func event_matches_hashtag(_ ev: NostrEvent, hashtags: [String]) -> Bool {
|
||||
@@ -106,33 +102,3 @@ func event_matches_filter(_ ev: NostrEvent, filter: NostrFilter) -> Bool {
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func handle_subid_event(pool: RelayPool, relay_id: RelayURL, ev: NostrConnectionEvent, handle: (String, NostrEvent) -> ()) -> (String?, Bool) {
|
||||
switch ev {
|
||||
case .ws_connection_event:
|
||||
return (nil, false)
|
||||
|
||||
case .nostr_event(let res):
|
||||
switch res {
|
||||
case .event(let ev_subid, let ev):
|
||||
handle(ev_subid, ev)
|
||||
return (ev_subid, false)
|
||||
|
||||
case .ok:
|
||||
return (nil, false)
|
||||
|
||||
case .notice(let note):
|
||||
if note.contains("Too many subscription filters") {
|
||||
// TODO: resend filters?
|
||||
pool.reconnect(to: [relay_id])
|
||||
}
|
||||
return (nil, false)
|
||||
|
||||
case .eose(let subid):
|
||||
return (subid, true)
|
||||
|
||||
case .auth:
|
||||
return (nil, false)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -14,6 +14,7 @@ struct SearchHomeView: View {
|
||||
@StateObject var model: SearchHomeModel
|
||||
@State var search: String = ""
|
||||
@FocusState private var isFocused: Bool
|
||||
@State var loadingTask: Task<Void, Never>?
|
||||
|
||||
func content_filter(_ fstate: FilterState) -> ((NostrEvent) -> Bool) {
|
||||
var filters = ContentFilters.defaults(damus_state: damus_state)
|
||||
@@ -84,8 +85,8 @@ struct SearchHomeView: View {
|
||||
)
|
||||
.refreshable {
|
||||
// Fetch new information by unsubscribing and resubscribing to the relay
|
||||
model.unsubscribe()
|
||||
model.subscribe()
|
||||
loadingTask?.cancel()
|
||||
loadingTask = Task { await model.load() }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -93,8 +94,8 @@ struct SearchHomeView: View {
|
||||
SearchResultsView(damus_state: damus_state, search: $search)
|
||||
.refreshable {
|
||||
// Fetch new information by unsubscribing and resubscribing to the relay
|
||||
model.unsubscribe()
|
||||
model.subscribe()
|
||||
loadingTask?.cancel()
|
||||
loadingTask = Task { await model.load() }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -129,11 +130,11 @@ struct SearchHomeView: View {
|
||||
}
|
||||
.onAppear {
|
||||
if model.events.events.isEmpty {
|
||||
model.subscribe()
|
||||
loadingTask = Task { await model.load() }
|
||||
}
|
||||
}
|
||||
.onDisappear {
|
||||
model.unsubscribe()
|
||||
loadingTask?.cancel()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -77,7 +77,8 @@ struct SearchingEventView: View {
|
||||
}
|
||||
|
||||
case .event(let note_id):
|
||||
find_event(state: state, query: .event(evid: note_id)) { res in
|
||||
Task {
|
||||
let res = await state.nostrNetwork.findEvent(query: .event(evid: note_id))
|
||||
guard case .event(let ev) = res else {
|
||||
self.search_state = .not_found
|
||||
return
|
||||
@@ -85,7 +86,8 @@ struct SearchingEventView: View {
|
||||
self.search_state = .found(ev)
|
||||
}
|
||||
case .profile(let pubkey):
|
||||
find_event(state: state, query: .profile(pubkey: pubkey)) { res in
|
||||
Task {
|
||||
let res = await state.nostrNetwork.findEvent(query: .profile(pubkey: pubkey))
|
||||
guard case .profile(let pubkey) = res else {
|
||||
self.search_state = .not_found
|
||||
return
|
||||
@@ -93,7 +95,8 @@ struct SearchingEventView: View {
|
||||
self.search_state = .found_profile(pubkey)
|
||||
}
|
||||
case .naddr(let naddr):
|
||||
naddrLookup(damus_state: state, naddr: naddr) { res in
|
||||
Task {
|
||||
let res = await state.nostrNetwork.lookup(naddr: naddr)
|
||||
guard let res = res else {
|
||||
self.search_state = .not_found
|
||||
return
|
||||
|
||||
Reference in New Issue
Block a user