This commit is contained in:
bas-vk 2016-04-01 12:54:26 +00:00
commit 8e67e36dba
4 changed files with 109 additions and 25 deletions

View file

@ -33,6 +33,8 @@ import (
"github.com/ethereum/go-ethereum/ethdb" "github.com/ethereum/go-ethereum/ethdb"
"github.com/ethereum/go-ethereum/event" "github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/rpc" "github.com/ethereum/go-ethereum/rpc"
"github.com/ethereum/go-ethereum/logger"
"github.com/ethereum/go-ethereum/logger/glog"
) )
var ( var (
@ -201,8 +203,9 @@ func (s *PublicFilterAPI) NewPendingTransactionFilter() (string, error) {
return externalId, nil return externalId, nil
} }
// newLogFilter creates a new log filter. // newLogFilter creates a new log filter. If emit is not nil logs are published
func (s *PublicFilterAPI) newLogFilter(earliest, latest int64, addresses []common.Address, topics [][]common.Hash) (int, error) { // over the channel, otherwise logs are queued and can be obtained on the next poll.
func (s *PublicFilterAPI) newLogFilter(earliest, latest int64, addresses []common.Address, topics [][]common.Hash, emit chan<- *event.Event) (int, error) {
s.logMu.Lock() s.logMu.Lock()
defer s.logMu.Unlock() defer s.logMu.Unlock()
@ -219,11 +222,14 @@ func (s *PublicFilterAPI) newLogFilter(earliest, latest int64, addresses []commo
filter.SetAddresses(addresses) filter.SetAddresses(addresses)
filter.SetTopics(topics) filter.SetTopics(topics)
filter.LogCallback = func(log *vm.Log, removed bool) { filter.LogCallback = func(log *vm.Log, removed bool) {
s.logMu.Lock() if emit == nil { // queue logs for next poll
defer s.logMu.Unlock() s.logMu.Lock()
defer s.logMu.Unlock()
if queue := s.logQueue[id]; queue != nil { if queue := s.logQueue[id]; queue != nil {
queue.add(vmlog{log, removed}) queue.add(vmlog{log, removed})
}
} else {
emit <- &event.Event{Time: time.Now(), Data: vmlog{log, removed}}
} }
} }
@ -364,10 +370,11 @@ func (s *PublicFilterAPI) NewFilter(args NewFilterArgs) (string, error) {
var id int var id int
if len(args.Addresses) > 0 { if len(args.Addresses) > 0 {
id, err = s.newLogFilter(args.FromBlock.Int64(), args.ToBlock.Int64(), args.Addresses, args.Topics) id, err = s.newLogFilter(args.FromBlock.Int64(), args.ToBlock.Int64(), args.Addresses, args.Topics, nil)
} else { } else {
id, err = s.newLogFilter(args.FromBlock.Int64(), args.ToBlock.Int64(), nil, args.Topics) id, err = s.newLogFilter(args.FromBlock.Int64(), args.ToBlock.Int64(), nil, args.Topics, nil)
} }
if err != nil { if err != nil {
return "", err return "", err
} }
@ -379,6 +386,64 @@ func (s *PublicFilterAPI) NewFilter(args NewFilterArgs) (string, error) {
return externalId, nil return externalId, nil
} }
// logsSubscription implements the event.Event interface and wraps a log subscription.
type logsSubscription struct {
api *PublicFilterAPI
filterId string
logs chan *event.Event
}
// Chan returns the channel where logs are published on for this subscription.
func (sub *logsSubscription) Chan() <-chan *event.Event {
return sub.logs
}
// Unsubscribe will uninstall the filter and closes the emit channel
func (sub *logsSubscription) Unsubscribe() {
if sub.api.UninstallFilter(sub.filterId) {
close(sub.logs)
}
}
// Logs creates a subscription and emits logs that match the given criteria.
func (s *PublicFilterAPI) Logs(args NewFilterArgs) (rpc.Subscription, error) {
externalId, err := newFilterId()
if err != nil {
return rpc.Subscription{}, err
}
logs := make(chan *event.Event)
// don't use earliest and latest since these are not applicable for subscriptions
var id int
if len(args.Addresses) > 0 {
id, err = s.newLogFilter(-1, -1, args.Addresses, args.Topics, logs)
} else {
id, err = s.newLogFilter(-1, -1, nil, args.Topics, logs)
}
if err != nil {
return rpc.Subscription{}, err
}
s.filterMapMu.Lock()
s.filterMapping[externalId] = id
s.filterMapMu.Unlock()
sub := &logsSubscription{s, externalId, logs}
format := func(rawLog interface{}) interface{} {
if log, ok := rawLog.(vmlog); ok {
rpcLogs := toRPCLogs(vm.Logs{log.Log}, log.Removed)
return &rpcLogs[0]
}
glog.V(logger.Warn).Infof("unexpected log type %T\n", rawLog)
return nil
}
return rpc.NewSubscriptionWithOutputFormat(sub, format), nil
}
// GetLogs returns the logs matching the given argument. // GetLogs returns the logs matching the given argument.
func (s *PublicFilterAPI) GetLogs(args NewFilterArgs) []vmlog { func (s *PublicFilterAPI) GetLogs(args NewFilterArgs) []vmlog {
filter := New(s.chainDb) filter := New(s.chainDb)
@ -514,6 +579,22 @@ type vmlog struct {
Removed bool `json:"removed"` Removed bool `json:"removed"`
} }
func (l *vmlog) MarshalJSON() ([]byte, error) {
fields := map[string]interface{}{
"address": l.Address,
"data": fmt.Sprintf("%#x", l.Data),
"blockNumber": fmt.Sprintf("%#x", l.BlockNumber),
"logIndex": fmt.Sprintf("%#x", l.Index),
"blockHash": l.BlockHash,
"transactionHash": l.TxHash,
"transactionIndex": fmt.Sprintf("%#x", l.TxIndex),
"topics": l.Topics,
"removed": l.Removed,
}
return json.Marshal(fields)
}
type logQueue struct { type logQueue struct {
mu sync.Mutex mu sync.Mutex
@ -564,9 +645,10 @@ func (l *hashQueue) get() []common.Hash {
return tmp return tmp
} }
// newFilterId generates a new random filter identifier that can be exposed to the outer world. By publishing random // newFilterId generates a new random filter identifier that can be exposed to
// identifiers it is not feasible for DApp's to guess filter id's for other DApp's and uninstall or poll for them // the world. By publishing random identifiers it is not feasible for DApp's
// causing the affected DApp to miss data. // to guess filter id's for other DApp's and uninstall or poll for them causing
// the targeted DApp to miss data.
func newFilterId() (string, error) { func newFilterId() (string, error) {
var subid [16]byte var subid [16]byte
n, _ := rand.Read(subid[:]) n, _ := rand.Read(subid[:])

View file

@ -164,7 +164,7 @@ func (fs *FilterSystem) filterLoop() {
fs.filterMu.RLock() fs.filterMu.RLock()
for _, filter := range fs.logFilters { for _, filter := range fs.logFilters {
if filter.LogCallback != nil && !filter.created.After(event.Time) { if filter.LogCallback != nil && !filter.created.After(event.Time) {
for _, removedLog := range ev.Logs { for _, removedLog := range filter.FilterLogs(ev.Logs) {
filter.LogCallback(removedLog, true) filter.LogCallback(removedLog, true)
} }
} }

View file

@ -225,14 +225,14 @@ func (s *Server) Stop() {
} }
} }
// sendNotification will create a notification from the given event by serializing member fields of the event. // sendNotification will send a notification to the client when the given event
// It will then send the notification to the client, when it fails the codec is closed. When the event has multiple // is not nil. It will close the codec when the subscription could not be send.
// fields an array of values is returned.
func sendNotification(codec ServerCodec, subid string, event interface{}) { func sendNotification(codec ServerCodec, subid string, event interface{}) {
notification := codec.CreateNotification(subid, event) if event != nil {
notification := codec.CreateNotification(subid, event)
if err := codec.Write(notification); err != nil { if err := codec.Write(notification); err != nil {
codec.Close() codec.Close()
}
} }
} }

View file

@ -136,9 +136,9 @@ func defaultSubscriptionOutputFormatter(data interface{}) interface{} {
// Subscription is used by the server to send notifications to the client // Subscription is used by the server to send notifications to the client
type Subscription struct { type Subscription struct {
sub event.Subscription sub event.Subscription // event generator
match SubscriptionMatcher match SubscriptionMatcher // allows for event filtering
format SubscriptionOutputFormat format SubscriptionOutputFormat // allows for event output formatting (e.g. add custom fields)
} }
// NewSubscription create a new RPC subscription // NewSubscription create a new RPC subscription
@ -146,13 +146,15 @@ func NewSubscription(sub event.Subscription) Subscription {
return Subscription{sub, nil, defaultSubscriptionOutputFormatter} return Subscription{sub, nil, defaultSubscriptionOutputFormatter}
} }
// NewSubscriptionWithOutputFormat create a new RPC subscription which a custom notification output format // NewSubscriptionWithOutputFormat creates a new RPC subscription which a custom notification output format.
// If formatter returns nil the event is discarded.
func NewSubscriptionWithOutputFormat(sub event.Subscription, formatter SubscriptionOutputFormat) Subscription { func NewSubscriptionWithOutputFormat(sub event.Subscription, formatter SubscriptionOutputFormat) Subscription {
return Subscription{sub, nil, formatter} return Subscription{sub, nil, formatter}
} }
// NewSubscriptionFiltered will create a new subscription. For each raised event the given matcher is // NewSubscriptionFiltered will create a new subscription. For each raised event the given matcher is
// called. If it returns true the event is send as notification to the client, otherwise it is ignored. // called. If it returns true the event is send to the client, otherwise it is discarded. This can be
// used in situations where custom filtering is necessary.
func NewSubscriptionFiltered(sub event.Subscription, match SubscriptionMatcher) Subscription { func NewSubscriptionFiltered(sub event.Subscription, match SubscriptionMatcher) Subscription {
return Subscription{sub, match, defaultSubscriptionOutputFormatter} return Subscription{sub, match, defaultSubscriptionOutputFormatter}
} }