Improve network database querying
This commit is contained in:
@@ -1,6 +1,7 @@
|
|||||||
package network
|
package network
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
@@ -115,13 +116,17 @@ func (s *StorageInterface) Get(key string) (record.Record, error) {
|
|||||||
// Query returns a an iterator for the supplied query.
|
// Query returns a an iterator for the supplied query.
|
||||||
func (s *StorageInterface) Query(q *query.Query, local, internal bool) (*iterator.Iterator, error) {
|
func (s *StorageInterface) Query(q *query.Query, local, internal bool) (*iterator.Iterator, error) {
|
||||||
it := iterator.New()
|
it := iterator.New()
|
||||||
go s.processQuery(q, it)
|
|
||||||
// TODO: check local and internal
|
module.StartWorker("connection query", func(_ context.Context) error {
|
||||||
|
s.processQuery(q, it)
|
||||||
|
return nil
|
||||||
|
})
|
||||||
|
|
||||||
return it, nil
|
return it, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *StorageInterface) processQuery(q *query.Query, it *iterator.Iterator) {
|
func (s *StorageInterface) processQuery(q *query.Query, it *iterator.Iterator) {
|
||||||
|
var matches bool
|
||||||
pid, scope, _, ok := parseDBKey(q.DatabaseKeyPrefix())
|
pid, scope, _, ok := parseDBKey(q.DatabaseKeyPrefix())
|
||||||
if !ok {
|
if !ok {
|
||||||
it.Finish(nil)
|
it.Finish(nil)
|
||||||
@@ -131,33 +136,42 @@ func (s *StorageInterface) processQuery(q *query.Query, it *iterator.Iterator) {
|
|||||||
if pid == process.UndefinedProcessID {
|
if pid == process.UndefinedProcessID {
|
||||||
// processes
|
// processes
|
||||||
for _, proc := range process.All() {
|
for _, proc := range process.All() {
|
||||||
proc.Lock()
|
func() {
|
||||||
if q.Matches(proc) {
|
proc.Lock()
|
||||||
|
defer proc.Unlock()
|
||||||
|
matches = q.Matches(proc)
|
||||||
|
}()
|
||||||
|
if matches {
|
||||||
it.Next <- proc
|
it.Next <- proc
|
||||||
}
|
}
|
||||||
proc.Unlock()
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if scope == "" || scope == "dns" {
|
if scope == "" || scope == "dns" {
|
||||||
// dns scopes only
|
// dns scopes only
|
||||||
for _, dnsConn := range dnsConns.clone() {
|
for _, dnsConn := range dnsConns.clone() {
|
||||||
dnsConn.Lock()
|
func() {
|
||||||
if q.Matches(dnsConn) {
|
dnsConn.Lock()
|
||||||
|
defer dnsConn.Unlock()
|
||||||
|
matches = q.Matches(dnsConn)
|
||||||
|
}()
|
||||||
|
if matches {
|
||||||
it.Next <- dnsConn
|
it.Next <- dnsConn
|
||||||
}
|
}
|
||||||
dnsConn.Unlock()
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if scope == "" || scope == "ip" {
|
if scope == "" || scope == "ip" {
|
||||||
// connections
|
// connections
|
||||||
for _, conn := range conns.clone() {
|
for _, conn := range conns.clone() {
|
||||||
conn.Lock()
|
func() {
|
||||||
if q.Matches(conn) {
|
conn.Lock()
|
||||||
|
defer conn.Unlock()
|
||||||
|
matches = q.Matches(conn)
|
||||||
|
}()
|
||||||
|
if matches {
|
||||||
it.Next <- conn
|
it.Next <- conn
|
||||||
}
|
}
|
||||||
conn.Unlock()
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user