This repository was archived by the owner on Nov 15, 2023. It is now read-only.

Storage changes subscription - #464

Merged
svyatonik merged 10 commits into
masterfrom
td-storage-events
Aug 1, 2018
Merged

Storage changes subscription#464
svyatonik merged 10 commits into
masterfrom
td-storage-events

Conversation

@tomusdrw

Copy link
Copy Markdown
Contributor
  • Exposing storage changes via RPC pub-sub
  • Allows one to subscribe to particular storage keys

CC @jacogr

@tomusdrwtomusdrw added A0-please_review Pull request needs code review. M6-rpcapi labels Jul 31, 2018
Comment threadsubstrate/rpc/src/chain/mod.rs Outdated
#[pubsub(name = "chain_newHead")] {
/// New head subscription
#[rpc(name = "subscribe_newHead")]
#[rpc(name = "subscribe_newHead", alias = ["chain_subscribeNewHead", ])]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Great. Maybe it is better to have the "old" as the alias and the "new" one as the name?

(e.g. chain_subscribeNewHead is actually probably the preferred one to match since it aligns with other RPCs, my gut tells me the preferred one should be the "default")

@dvdplmdvdplm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good stuff. I can't say I understand it all, but overall very readable and interesting.

Comment threadsubstrate/client/src/client.rs Outdated

/// Get storage changes event stream.
///
/// Passing `None` as keys subscribes to all possible keys

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

typo: …as keys should be …as key

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rephrased the whole sentence, hope it's clearer now.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some(storage_update) = storage_update {
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't we emit a re-org event with some meta data about what changed and let interested subscribers re-fetch?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, that's possible, although it's not easy to get storage changes for blocks that are already imported/executed. Re-fetching changes would mean to re-execute the blocks, but I suppose based on the filter_keys we could just return all the storage values in the re-orged blocks.

subscribers.extend(listeners.iter());
}

if has_wildcard || listeners.is_some() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't you check for !susbscribers.is_empty() here? Or is it faster to do it this way?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm actually interested only in subscribers for that particular key. So if there is a set of changes:

[(1, Some(2)), (2, None), (3, Some(4))]

but we have no wildcard_listeners and only listener for key=1, changes vector will only contain [(StorageKey(1), Some(StorageData(2))]

filter: filter.clone(),
})).is_err()
},
None => false,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If I read this right, if the subscriber is gone when we get here, then we assume they've been removed properly already, so we're returning false to avoid calling remove_subscriber() again for them? It's a bit unclear to me how they can still be in the subscribers collection though, can you elaborate on that?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Indeed, I could actually .expect() here, since if the structure is consistent the subscribers should always be in self.sinks. The check here is superfluous.
I could refactor to:

let&(ref sink,ref filter) = self.sinks.get(&subscriber).expect("subscribers returned from self.listeners are always in self.sinks; qed");let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}

or

ifletSome(&(ref sink,ref filter)) = matchself.sinks.get(&subscriber){
let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}}

Which one do you prefer?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Oh, actually I can't since .get() borrows immutably, so remove_subscriber has to be outside of the scope.

assert_eq!(notifications.listeners.len(), 2);
assert_eq!(notifications.wildcard_listeners.len(), 1);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The channels are closed here, correct?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes and since the receiving end is dropped, sending to such channel will trigger an error.

sink
.sink_map_err(|e| warn!("Error sending notifications: {:?}", e))
.send_all(stream)
// we ignore the resulting Stream (if the first stream is over we are unsubscribed)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

…is over…? Do you mean …is closed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

is over as in is finished/is done. Which means that the stream will not emit any more items.

/// Drain committed changes to an iterator.
///
/// Panics:
/// Will panic if there are any uncommitted prospective changes.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is "prospective" sort of like "pending"?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, those are changes that can still be easily discarded. You can see example of usage inside block builder:

  1. We run a transaction
  2. It produces a set of prospective changes
  3. If we detect that it's somehow invalid we discard the prospective changes
  4. If we accept the transaction we commit prospective changes.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?
self.storage_notifications.lock()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are you sure that it should be called before transaction is committed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. Moved the notification after commit and also guarded by the same if as block import notification.

@gavofyorkgavofyork left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Aside from the minor comment

@@ -0,0 +1,267 @@
// Copyright 2017 Parity Technologies (UK) Ltd.
// This file is part of Polkadot.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Substrate, not Polkadot :)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed.

@gavofyorkgavofyork removed the A0-please_review Pull request needs code review. label Aug 1, 2018
@svyatonik
svyatonik merged commit 757a721 into masterAug 1, 2018
@svyatonik
svyatonik deleted the td-storage-events branch August 1, 2018 12:29
dvdplm added a commit that referenced this pull request Aug 1, 2018
* master:
Collator for the "adder" (formerly basic-add) parachain and various small fixes (#438)
Storage changes subscription (#464)
Wasm execution optimizations (#466)
Fix the --key generation (#475)
Fix typo in service.rs (#472)
Fix session phase in early-exit (#453)
Make ping unidirectional (#458)
Update README.adoc
gavofyork pushed a commit that referenced this pull request Aug 10, 2018
* Initial implementation of storage events.
* Attaching storage events.
* Expose storage modification stream over RPC.
* Use FNV for hashing small keys.
* Fix and add tests.
* Swap alias and RPC name.
* Fix demo.
* Addressing review grumbles.
* Fix comment.
liuchengxu pushed a commit to autonomys/substrate that referenced this pull request Jun 3, 2022
helin6 pushed a commit to boolnetwork/substrate that referenced this pull request Jul 25, 2023
* Bump release version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update dependency version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Modify changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Move changelog entries from added to changed
Sign up for freeto subscribe to this conversation on GitHub. Already have an account? Sign in.

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@tomusdrw@dvdplm@gavofyork@svyatonik@jacogr
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
This repository was archived by the owner on Nov 15, 2023. It is now read-only.

Storage changes subscription - #464

Merged
svyatonik merged 10 commits into
masterfrom
td-storage-events
Aug 1, 2018
Merged

Storage changes subscription#464
svyatonik merged 10 commits into
masterfrom
td-storage-events

Conversation

@tomusdrw

Copy link
Copy Markdown
Contributor
  • Exposing storage changes via RPC pub-sub
  • Allows one to subscribe to particular storage keys

CC @jacogr

@tomusdrwtomusdrw added A0-please_review Pull request needs code review. M6-rpcapi labels Jul 31, 2018
Comment threadsubstrate/rpc/src/chain/mod.rs Outdated
#[pubsub(name = "chain_newHead")] {
/// New head subscription
#[rpc(name = "subscribe_newHead")]
#[rpc(name = "subscribe_newHead", alias = ["chain_subscribeNewHead", ])]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Great. Maybe it is better to have the "old" as the alias and the "new" one as the name?

(e.g. chain_subscribeNewHead is actually probably the preferred one to match since it aligns with other RPCs, my gut tells me the preferred one should be the "default")

@dvdplmdvdplm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good stuff. I can't say I understand it all, but overall very readable and interesting.

Comment threadsubstrate/client/src/client.rs Outdated

/// Get storage changes event stream.
///
/// Passing `None` as keys subscribes to all possible keys

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

typo: …as keys should be …as key

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rephrased the whole sentence, hope it's clearer now.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some(storage_update) = storage_update {
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't we emit a re-org event with some meta data about what changed and let interested subscribers re-fetch?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, that's possible, although it's not easy to get storage changes for blocks that are already imported/executed. Re-fetching changes would mean to re-execute the blocks, but I suppose based on the filter_keys we could just return all the storage values in the re-orged blocks.

subscribers.extend(listeners.iter());
}

if has_wildcard || listeners.is_some() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't you check for !susbscribers.is_empty() here? Or is it faster to do it this way?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm actually interested only in subscribers for that particular key. So if there is a set of changes:

[(1, Some(2)), (2, None), (3, Some(4))]

but we have no wildcard_listeners and only listener for key=1, changes vector will only contain [(StorageKey(1), Some(StorageData(2))]

filter: filter.clone(),
})).is_err()
},
None => false,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If I read this right, if the subscriber is gone when we get here, then we assume they've been removed properly already, so we're returning false to avoid calling remove_subscriber() again for them? It's a bit unclear to me how they can still be in the subscribers collection though, can you elaborate on that?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Indeed, I could actually .expect() here, since if the structure is consistent the subscribers should always be in self.sinks. The check here is superfluous.
I could refactor to:

let&(ref sink,ref filter) = self.sinks.get(&subscriber).expect("subscribers returned from self.listeners are always in self.sinks; qed");let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}

or

ifletSome(&(ref sink,ref filter)) = matchself.sinks.get(&subscriber){
let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}}

Which one do you prefer?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Oh, actually I can't since .get() borrows immutably, so remove_subscriber has to be outside of the scope.

assert_eq!(notifications.listeners.len(), 2);
assert_eq!(notifications.wildcard_listeners.len(), 1);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The channels are closed here, correct?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes and since the receiving end is dropped, sending to such channel will trigger an error.

sink
.sink_map_err(|e| warn!("Error sending notifications: {:?}", e))
.send_all(stream)
// we ignore the resulting Stream (if the first stream is over we are unsubscribed)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

…is over…? Do you mean …is closed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

is over as in is finished/is done. Which means that the stream will not emit any more items.

/// Drain committed changes to an iterator.
///
/// Panics:
/// Will panic if there are any uncommitted prospective changes.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is "prospective" sort of like "pending"?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, those are changes that can still be easily discarded. You can see example of usage inside block builder:

  1. We run a transaction
  2. It produces a set of prospective changes
  3. If we detect that it's somehow invalid we discard the prospective changes
  4. If we accept the transaction we commit prospective changes.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?
self.storage_notifications.lock()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are you sure that it should be called before transaction is committed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. Moved the notification after commit and also guarded by the same if as block import notification.

@gavofyorkgavofyork left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Aside from the minor comment

@@ -0,0 +1,267 @@
// Copyright 2017 Parity Technologies (UK) Ltd.
// This file is part of Polkadot.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Substrate, not Polkadot :)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed.

@gavofyorkgavofyork removed the A0-please_review Pull request needs code review. label Aug 1, 2018
@svyatonik
svyatonik merged commit 757a721 into masterAug 1, 2018
@svyatonik
svyatonik deleted the td-storage-events branch August 1, 2018 12:29
dvdplm added a commit that referenced this pull request Aug 1, 2018
* master:
Collator for the "adder" (formerly basic-add) parachain and various small fixes (#438)
Storage changes subscription (#464)
Wasm execution optimizations (#466)
Fix the --key generation (#475)
Fix typo in service.rs (#472)
Fix session phase in early-exit (#453)
Make ping unidirectional (#458)
Update README.adoc
gavofyork pushed a commit that referenced this pull request Aug 10, 2018
* Initial implementation of storage events.
* Attaching storage events.
* Expose storage modification stream over RPC.
* Use FNV for hashing small keys.
* Fix and add tests.
* Swap alias and RPC name.
* Fix demo.
* Addressing review grumbles.
* Fix comment.
liuchengxu pushed a commit to autonomys/substrate that referenced this pull request Jun 3, 2022
helin6 pushed a commit to boolnetwork/substrate that referenced this pull request Jul 25, 2023
* Bump release version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update dependency version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Modify changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Move changelog entries from added to changed
Sign up for freeto subscribe to this conversation on GitHub. Already have an account? Sign in.

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@tomusdrw@dvdplm@gavofyork@svyatonik@jacogr
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
This repository was archived by the owner on Nov 15, 2023. It is now read-only.

Storage changes subscription - #464

Merged
svyatonik merged 10 commits into
masterfrom
td-storage-events
Aug 1, 2018
Merged

Storage changes subscription#464
svyatonik merged 10 commits into
masterfrom
td-storage-events

Conversation

@tomusdrw

Copy link
Copy Markdown
Contributor
  • Exposing storage changes via RPC pub-sub
  • Allows one to subscribe to particular storage keys

CC @jacogr

@tomusdrwtomusdrw added A0-please_review Pull request needs code review. M6-rpcapi labels Jul 31, 2018
Comment threadsubstrate/rpc/src/chain/mod.rs Outdated
#[pubsub(name = "chain_newHead")] {
/// New head subscription
#[rpc(name = "subscribe_newHead")]
#[rpc(name = "subscribe_newHead", alias = ["chain_subscribeNewHead", ])]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Great. Maybe it is better to have the "old" as the alias and the "new" one as the name?

(e.g. chain_subscribeNewHead is actually probably the preferred one to match since it aligns with other RPCs, my gut tells me the preferred one should be the "default")

@dvdplmdvdplm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good stuff. I can't say I understand it all, but overall very readable and interesting.

Comment threadsubstrate/client/src/client.rs Outdated

/// Get storage changes event stream.
///
/// Passing `None` as keys subscribes to all possible keys

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

typo: …as keys should be …as key

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rephrased the whole sentence, hope it's clearer now.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some(storage_update) = storage_update {
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't we emit a re-org event with some meta data about what changed and let interested subscribers re-fetch?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, that's possible, although it's not easy to get storage changes for blocks that are already imported/executed. Re-fetching changes would mean to re-execute the blocks, but I suppose based on the filter_keys we could just return all the storage values in the re-orged blocks.

subscribers.extend(listeners.iter());
}

if has_wildcard || listeners.is_some() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't you check for !susbscribers.is_empty() here? Or is it faster to do it this way?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm actually interested only in subscribers for that particular key. So if there is a set of changes:

[(1, Some(2)), (2, None), (3, Some(4))]

but we have no wildcard_listeners and only listener for key=1, changes vector will only contain [(StorageKey(1), Some(StorageData(2))]

filter: filter.clone(),
})).is_err()
},
None => false,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If I read this right, if the subscriber is gone when we get here, then we assume they've been removed properly already, so we're returning false to avoid calling remove_subscriber() again for them? It's a bit unclear to me how they can still be in the subscribers collection though, can you elaborate on that?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Indeed, I could actually .expect() here, since if the structure is consistent the subscribers should always be in self.sinks. The check here is superfluous.
I could refactor to:

let&(ref sink,ref filter) = self.sinks.get(&subscriber).expect("subscribers returned from self.listeners are always in self.sinks; qed");let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}

or

ifletSome(&(ref sink,ref filter)) = matchself.sinks.get(&subscriber){
let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}}

Which one do you prefer?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Oh, actually I can't since .get() borrows immutably, so remove_subscriber has to be outside of the scope.

assert_eq!(notifications.listeners.len(), 2);
assert_eq!(notifications.wildcard_listeners.len(), 1);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The channels are closed here, correct?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes and since the receiving end is dropped, sending to such channel will trigger an error.

sink
.sink_map_err(|e| warn!("Error sending notifications: {:?}", e))
.send_all(stream)
// we ignore the resulting Stream (if the first stream is over we are unsubscribed)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

…is over…? Do you mean …is closed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

is over as in is finished/is done. Which means that the stream will not emit any more items.

/// Drain committed changes to an iterator.
///
/// Panics:
/// Will panic if there are any uncommitted prospective changes.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is "prospective" sort of like "pending"?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, those are changes that can still be easily discarded. You can see example of usage inside block builder:

  1. We run a transaction
  2. It produces a set of prospective changes
  3. If we detect that it's somehow invalid we discard the prospective changes
  4. If we accept the transaction we commit prospective changes.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?
self.storage_notifications.lock()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are you sure that it should be called before transaction is committed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. Moved the notification after commit and also guarded by the same if as block import notification.

@gavofyorkgavofyork left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Aside from the minor comment

@@ -0,0 +1,267 @@
// Copyright 2017 Parity Technologies (UK) Ltd.
// This file is part of Polkadot.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Substrate, not Polkadot :)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed.

@gavofyorkgavofyork removed the A0-please_review Pull request needs code review. label Aug 1, 2018
@svyatonik
svyatonik merged commit 757a721 into masterAug 1, 2018
@svyatonik
svyatonik deleted the td-storage-events branch August 1, 2018 12:29
dvdplm added a commit that referenced this pull request Aug 1, 2018
* master:
Collator for the "adder" (formerly basic-add) parachain and various small fixes (#438)
Storage changes subscription (#464)
Wasm execution optimizations (#466)
Fix the --key generation (#475)
Fix typo in service.rs (#472)
Fix session phase in early-exit (#453)
Make ping unidirectional (#458)
Update README.adoc
gavofyork pushed a commit that referenced this pull request Aug 10, 2018
* Initial implementation of storage events.
* Attaching storage events.
* Expose storage modification stream over RPC.
* Use FNV for hashing small keys.
* Fix and add tests.
* Swap alias and RPC name.
* Fix demo.
* Addressing review grumbles.
* Fix comment.
liuchengxu pushed a commit to autonomys/substrate that referenced this pull request Jun 3, 2022
helin6 pushed a commit to boolnetwork/substrate that referenced this pull request Jul 25, 2023
* Bump release version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update dependency version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Modify changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Move changelog entries from added to changed
Sign up for freeto subscribe to this conversation on GitHub. Already have an account? Sign in.

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@tomusdrw@dvdplm@gavofyork@svyatonik@jacogr
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
This repository was archived by the owner on Nov 15, 2023. It is now read-only.

Storage changes subscription - #464

Merged
svyatonik merged 10 commits into
masterfrom
td-storage-events
Aug 1, 2018
Merged

Storage changes subscription#464
svyatonik merged 10 commits into
masterfrom
td-storage-events

Conversation

@tomusdrw

Copy link
Copy Markdown
Contributor
  • Exposing storage changes via RPC pub-sub
  • Allows one to subscribe to particular storage keys

CC @jacogr

@tomusdrwtomusdrw added A0-please_review Pull request needs code review. M6-rpcapi labels Jul 31, 2018
Comment threadsubstrate/rpc/src/chain/mod.rs Outdated
#[pubsub(name = "chain_newHead")] {
/// New head subscription
#[rpc(name = "subscribe_newHead")]
#[rpc(name = "subscribe_newHead", alias = ["chain_subscribeNewHead", ])]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Great. Maybe it is better to have the "old" as the alias and the "new" one as the name?

(e.g. chain_subscribeNewHead is actually probably the preferred one to match since it aligns with other RPCs, my gut tells me the preferred one should be the "default")

@dvdplmdvdplm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good stuff. I can't say I understand it all, but overall very readable and interesting.

Comment threadsubstrate/client/src/client.rs Outdated

/// Get storage changes event stream.
///
/// Passing `None` as keys subscribes to all possible keys

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

typo: …as keys should be …as key

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rephrased the whole sentence, hope it's clearer now.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some(storage_update) = storage_update {
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't we emit a re-org event with some meta data about what changed and let interested subscribers re-fetch?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, that's possible, although it's not easy to get storage changes for blocks that are already imported/executed. Re-fetching changes would mean to re-execute the blocks, but I suppose based on the filter_keys we could just return all the storage values in the re-orged blocks.

subscribers.extend(listeners.iter());
}

if has_wildcard || listeners.is_some() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't you check for !susbscribers.is_empty() here? Or is it faster to do it this way?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm actually interested only in subscribers for that particular key. So if there is a set of changes:

[(1, Some(2)), (2, None), (3, Some(4))]

but we have no wildcard_listeners and only listener for key=1, changes vector will only contain [(StorageKey(1), Some(StorageData(2))]

filter: filter.clone(),
})).is_err()
},
None => false,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If I read this right, if the subscriber is gone when we get here, then we assume they've been removed properly already, so we're returning false to avoid calling remove_subscriber() again for them? It's a bit unclear to me how they can still be in the subscribers collection though, can you elaborate on that?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Indeed, I could actually .expect() here, since if the structure is consistent the subscribers should always be in self.sinks. The check here is superfluous.
I could refactor to:

let&(ref sink,ref filter) = self.sinks.get(&subscriber).expect("subscribers returned from self.listeners are always in self.sinks; qed");let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}

or

ifletSome(&(ref sink,ref filter)) = matchself.sinks.get(&subscriber){
let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}}

Which one do you prefer?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Oh, actually I can't since .get() borrows immutably, so remove_subscriber has to be outside of the scope.

assert_eq!(notifications.listeners.len(), 2);
assert_eq!(notifications.wildcard_listeners.len(), 1);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The channels are closed here, correct?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes and since the receiving end is dropped, sending to such channel will trigger an error.

sink
.sink_map_err(|e| warn!("Error sending notifications: {:?}", e))
.send_all(stream)
// we ignore the resulting Stream (if the first stream is over we are unsubscribed)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

…is over…? Do you mean …is closed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

is over as in is finished/is done. Which means that the stream will not emit any more items.

/// Drain committed changes to an iterator.
///
/// Panics:
/// Will panic if there are any uncommitted prospective changes.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is "prospective" sort of like "pending"?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, those are changes that can still be easily discarded. You can see example of usage inside block builder:

  1. We run a transaction
  2. It produces a set of prospective changes
  3. If we detect that it's somehow invalid we discard the prospective changes
  4. If we accept the transaction we commit prospective changes.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?
self.storage_notifications.lock()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are you sure that it should be called before transaction is committed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. Moved the notification after commit and also guarded by the same if as block import notification.

@gavofyorkgavofyork left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Aside from the minor comment

@@ -0,0 +1,267 @@
// Copyright 2017 Parity Technologies (UK) Ltd.
// This file is part of Polkadot.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Substrate, not Polkadot :)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed.

@gavofyorkgavofyork removed the A0-please_review Pull request needs code review. label Aug 1, 2018
@svyatonik
svyatonik merged commit 757a721 into masterAug 1, 2018
@svyatonik
svyatonik deleted the td-storage-events branch August 1, 2018 12:29
dvdplm added a commit that referenced this pull request Aug 1, 2018
* master:
Collator for the "adder" (formerly basic-add) parachain and various small fixes (#438)
Storage changes subscription (#464)
Wasm execution optimizations (#466)
Fix the --key generation (#475)
Fix typo in service.rs (#472)
Fix session phase in early-exit (#453)
Make ping unidirectional (#458)
Update README.adoc
gavofyork pushed a commit that referenced this pull request Aug 10, 2018
* Initial implementation of storage events.
* Attaching storage events.
* Expose storage modification stream over RPC.
* Use FNV for hashing small keys.
* Fix and add tests.
* Swap alias and RPC name.
* Fix demo.
* Addressing review grumbles.
* Fix comment.
liuchengxu pushed a commit to autonomys/substrate that referenced this pull request Jun 3, 2022
helin6 pushed a commit to boolnetwork/substrate that referenced this pull request Jul 25, 2023
* Bump release version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update dependency version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Modify changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Move changelog entries from added to changed
Sign up for freeto subscribe to this conversation on GitHub. Already have an account? Sign in.

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@tomusdrw@dvdplm@gavofyork@svyatonik@jacogr
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
This repository was archived by the owner on Nov 15, 2023. It is now read-only.

Storage changes subscription - #464

Merged
svyatonik merged 10 commits into
masterfrom
td-storage-events
Aug 1, 2018
Merged

Storage changes subscription#464
svyatonik merged 10 commits into
masterfrom
td-storage-events

Conversation

@tomusdrw

Copy link
Copy Markdown
Contributor
  • Exposing storage changes via RPC pub-sub
  • Allows one to subscribe to particular storage keys

CC @jacogr

@tomusdrwtomusdrw added A0-please_review Pull request needs code review. M6-rpcapi labels Jul 31, 2018
Comment threadsubstrate/rpc/src/chain/mod.rs Outdated
#[pubsub(name = "chain_newHead")] {
/// New head subscription
#[rpc(name = "subscribe_newHead")]
#[rpc(name = "subscribe_newHead", alias = ["chain_subscribeNewHead", ])]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Great. Maybe it is better to have the "old" as the alias and the "new" one as the name?

(e.g. chain_subscribeNewHead is actually probably the preferred one to match since it aligns with other RPCs, my gut tells me the preferred one should be the "default")

@dvdplmdvdplm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good stuff. I can't say I understand it all, but overall very readable and interesting.

Comment threadsubstrate/client/src/client.rs Outdated

/// Get storage changes event stream.
///
/// Passing `None` as keys subscribes to all possible keys

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

typo: …as keys should be …as key

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rephrased the whole sentence, hope it's clearer now.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some(storage_update) = storage_update {
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't we emit a re-org event with some meta data about what changed and let interested subscribers re-fetch?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, that's possible, although it's not easy to get storage changes for blocks that are already imported/executed. Re-fetching changes would mean to re-execute the blocks, but I suppose based on the filter_keys we could just return all the storage values in the re-orged blocks.

subscribers.extend(listeners.iter());
}

if has_wildcard || listeners.is_some() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't you check for !susbscribers.is_empty() here? Or is it faster to do it this way?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm actually interested only in subscribers for that particular key. So if there is a set of changes:

[(1, Some(2)), (2, None), (3, Some(4))]

but we have no wildcard_listeners and only listener for key=1, changes vector will only contain [(StorageKey(1), Some(StorageData(2))]

filter: filter.clone(),
})).is_err()
},
None => false,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If I read this right, if the subscriber is gone when we get here, then we assume they've been removed properly already, so we're returning false to avoid calling remove_subscriber() again for them? It's a bit unclear to me how they can still be in the subscribers collection though, can you elaborate on that?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Indeed, I could actually .expect() here, since if the structure is consistent the subscribers should always be in self.sinks. The check here is superfluous.
I could refactor to:

let&(ref sink,ref filter) = self.sinks.get(&subscriber).expect("subscribers returned from self.listeners are always in self.sinks; qed");let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}

or

ifletSome(&(ref sink,ref filter)) = matchself.sinks.get(&subscriber){
let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}}

Which one do you prefer?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Oh, actually I can't since .get() borrows immutably, so remove_subscriber has to be outside of the scope.

assert_eq!(notifications.listeners.len(), 2);
assert_eq!(notifications.wildcard_listeners.len(), 1);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The channels are closed here, correct?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes and since the receiving end is dropped, sending to such channel will trigger an error.

sink
.sink_map_err(|e| warn!("Error sending notifications: {:?}", e))
.send_all(stream)
// we ignore the resulting Stream (if the first stream is over we are unsubscribed)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

…is over…? Do you mean …is closed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

is over as in is finished/is done. Which means that the stream will not emit any more items.

/// Drain committed changes to an iterator.
///
/// Panics:
/// Will panic if there are any uncommitted prospective changes.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is "prospective" sort of like "pending"?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, those are changes that can still be easily discarded. You can see example of usage inside block builder:

  1. We run a transaction
  2. It produces a set of prospective changes
  3. If we detect that it's somehow invalid we discard the prospective changes
  4. If we accept the transaction we commit prospective changes.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?
self.storage_notifications.lock()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are you sure that it should be called before transaction is committed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. Moved the notification after commit and also guarded by the same if as block import notification.

@gavofyorkgavofyork left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Aside from the minor comment

@@ -0,0 +1,267 @@
// Copyright 2017 Parity Technologies (UK) Ltd.
// This file is part of Polkadot.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Substrate, not Polkadot :)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed.

@gavofyorkgavofyork removed the A0-please_review Pull request needs code review. label Aug 1, 2018
@svyatonik
svyatonik merged commit 757a721 into masterAug 1, 2018
@svyatonik
svyatonik deleted the td-storage-events branch August 1, 2018 12:29
dvdplm added a commit that referenced this pull request Aug 1, 2018
* master:
Collator for the "adder" (formerly basic-add) parachain and various small fixes (#438)
Storage changes subscription (#464)
Wasm execution optimizations (#466)
Fix the --key generation (#475)
Fix typo in service.rs (#472)
Fix session phase in early-exit (#453)
Make ping unidirectional (#458)
Update README.adoc
gavofyork pushed a commit that referenced this pull request Aug 10, 2018
* Initial implementation of storage events.
* Attaching storage events.
* Expose storage modification stream over RPC.
* Use FNV for hashing small keys.
* Fix and add tests.
* Swap alias and RPC name.
* Fix demo.
* Addressing review grumbles.
* Fix comment.
liuchengxu pushed a commit to autonomys/substrate that referenced this pull request Jun 3, 2022
helin6 pushed a commit to boolnetwork/substrate that referenced this pull request Jul 25, 2023
* Bump release version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update dependency version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Modify changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Move changelog entries from added to changed
Sign up for freeto subscribe to this conversation on GitHub. Already have an account? Sign in.

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@tomusdrw@dvdplm@gavofyork@svyatonik@jacogr
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
This repository was archived by the owner on Nov 15, 2023. It is now read-only.

Storage changes subscription - #464

Merged
svyatonik merged 10 commits into
masterfrom
td-storage-events
Aug 1, 2018
Merged

Storage changes subscription#464
svyatonik merged 10 commits into
masterfrom
td-storage-events

Conversation

@tomusdrw

Copy link
Copy Markdown
Contributor
  • Exposing storage changes via RPC pub-sub
  • Allows one to subscribe to particular storage keys

CC @jacogr

@tomusdrwtomusdrw added A0-please_review Pull request needs code review. M6-rpcapi labels Jul 31, 2018
Comment threadsubstrate/rpc/src/chain/mod.rs Outdated
#[pubsub(name = "chain_newHead")] {
/// New head subscription
#[rpc(name = "subscribe_newHead")]
#[rpc(name = "subscribe_newHead", alias = ["chain_subscribeNewHead", ])]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Great. Maybe it is better to have the "old" as the alias and the "new" one as the name?

(e.g. chain_subscribeNewHead is actually probably the preferred one to match since it aligns with other RPCs, my gut tells me the preferred one should be the "default")

@dvdplmdvdplm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good stuff. I can't say I understand it all, but overall very readable and interesting.

Comment threadsubstrate/client/src/client.rs Outdated

/// Get storage changes event stream.
///
/// Passing `None` as keys subscribes to all possible keys

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

typo: …as keys should be …as key

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rephrased the whole sentence, hope it's clearer now.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some(storage_update) = storage_update {
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't we emit a re-org event with some meta data about what changed and let interested subscribers re-fetch?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, that's possible, although it's not easy to get storage changes for blocks that are already imported/executed. Re-fetching changes would mean to re-execute the blocks, but I suppose based on the filter_keys we could just return all the storage values in the re-orged blocks.

subscribers.extend(listeners.iter());
}

if has_wildcard || listeners.is_some() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't you check for !susbscribers.is_empty() here? Or is it faster to do it this way?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm actually interested only in subscribers for that particular key. So if there is a set of changes:

[(1, Some(2)), (2, None), (3, Some(4))]

but we have no wildcard_listeners and only listener for key=1, changes vector will only contain [(StorageKey(1), Some(StorageData(2))]

filter: filter.clone(),
})).is_err()
},
None => false,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If I read this right, if the subscriber is gone when we get here, then we assume they've been removed properly already, so we're returning false to avoid calling remove_subscriber() again for them? It's a bit unclear to me how they can still be in the subscribers collection though, can you elaborate on that?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Indeed, I could actually .expect() here, since if the structure is consistent the subscribers should always be in self.sinks. The check here is superfluous.
I could refactor to:

let&(ref sink,ref filter) = self.sinks.get(&subscriber).expect("subscribers returned from self.listeners are always in self.sinks; qed");let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}

or

ifletSome(&(ref sink,ref filter)) = matchself.sinks.get(&subscriber){
let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}}

Which one do you prefer?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Oh, actually I can't since .get() borrows immutably, so remove_subscriber has to be outside of the scope.

assert_eq!(notifications.listeners.len(), 2);
assert_eq!(notifications.wildcard_listeners.len(), 1);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The channels are closed here, correct?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes and since the receiving end is dropped, sending to such channel will trigger an error.

sink
.sink_map_err(|e| warn!("Error sending notifications: {:?}", e))
.send_all(stream)
// we ignore the resulting Stream (if the first stream is over we are unsubscribed)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

…is over…? Do you mean …is closed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

is over as in is finished/is done. Which means that the stream will not emit any more items.

/// Drain committed changes to an iterator.
///
/// Panics:
/// Will panic if there are any uncommitted prospective changes.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is "prospective" sort of like "pending"?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, those are changes that can still be easily discarded. You can see example of usage inside block builder:

  1. We run a transaction
  2. It produces a set of prospective changes
  3. If we detect that it's somehow invalid we discard the prospective changes
  4. If we accept the transaction we commit prospective changes.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?
self.storage_notifications.lock()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are you sure that it should be called before transaction is committed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. Moved the notification after commit and also guarded by the same if as block import notification.

@gavofyorkgavofyork left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Aside from the minor comment

@@ -0,0 +1,267 @@
// Copyright 2017 Parity Technologies (UK) Ltd.
// This file is part of Polkadot.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Substrate, not Polkadot :)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed.

@gavofyorkgavofyork removed the A0-please_review Pull request needs code review. label Aug 1, 2018
@svyatonik
svyatonik merged commit 757a721 into masterAug 1, 2018
@svyatonik
svyatonik deleted the td-storage-events branch August 1, 2018 12:29
dvdplm added a commit that referenced this pull request Aug 1, 2018
* master:
Collator for the "adder" (formerly basic-add) parachain and various small fixes (#438)
Storage changes subscription (#464)
Wasm execution optimizations (#466)
Fix the --key generation (#475)
Fix typo in service.rs (#472)
Fix session phase in early-exit (#453)
Make ping unidirectional (#458)
Update README.adoc
gavofyork pushed a commit that referenced this pull request Aug 10, 2018
* Initial implementation of storage events.
* Attaching storage events.
* Expose storage modification stream over RPC.
* Use FNV for hashing small keys.
* Fix and add tests.
* Swap alias and RPC name.
* Fix demo.
* Addressing review grumbles.
* Fix comment.
liuchengxu pushed a commit to autonomys/substrate that referenced this pull request Jun 3, 2022
helin6 pushed a commit to boolnetwork/substrate that referenced this pull request Jul 25, 2023
* Bump release version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update dependency version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Modify changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Move changelog entries from added to changed
Sign up for freeto subscribe to this conversation on GitHub. Already have an account? Sign in.

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@tomusdrw@dvdplm@gavofyork@svyatonik@jacogr
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
This repository was archived by the owner on Nov 15, 2023. It is now read-only.

Storage changes subscription - #464

Merged
svyatonik merged 10 commits into
masterfrom
td-storage-events
Aug 1, 2018
Merged

Storage changes subscription#464
svyatonik merged 10 commits into
masterfrom
td-storage-events

Conversation

@tomusdrw

Copy link
Copy Markdown
Contributor
  • Exposing storage changes via RPC pub-sub
  • Allows one to subscribe to particular storage keys

CC @jacogr

@tomusdrwtomusdrw added A0-please_review Pull request needs code review. M6-rpcapi labels Jul 31, 2018
Comment threadsubstrate/rpc/src/chain/mod.rs Outdated
#[pubsub(name = "chain_newHead")] {
/// New head subscription
#[rpc(name = "subscribe_newHead")]
#[rpc(name = "subscribe_newHead", alias = ["chain_subscribeNewHead", ])]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Great. Maybe it is better to have the "old" as the alias and the "new" one as the name?

(e.g. chain_subscribeNewHead is actually probably the preferred one to match since it aligns with other RPCs, my gut tells me the preferred one should be the "default")

@dvdplmdvdplm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good stuff. I can't say I understand it all, but overall very readable and interesting.

Comment threadsubstrate/client/src/client.rs Outdated

/// Get storage changes event stream.
///
/// Passing `None` as keys subscribes to all possible keys

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

typo: …as keys should be …as key

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rephrased the whole sentence, hope it's clearer now.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some(storage_update) = storage_update {
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't we emit a re-org event with some meta data about what changed and let interested subscribers re-fetch?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, that's possible, although it's not easy to get storage changes for blocks that are already imported/executed. Re-fetching changes would mean to re-execute the blocks, but I suppose based on the filter_keys we could just return all the storage values in the re-orged blocks.

subscribers.extend(listeners.iter());
}

if has_wildcard || listeners.is_some() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't you check for !susbscribers.is_empty() here? Or is it faster to do it this way?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm actually interested only in subscribers for that particular key. So if there is a set of changes:

[(1, Some(2)), (2, None), (3, Some(4))]

but we have no wildcard_listeners and only listener for key=1, changes vector will only contain [(StorageKey(1), Some(StorageData(2))]

filter: filter.clone(),
})).is_err()
},
None => false,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If I read this right, if the subscriber is gone when we get here, then we assume they've been removed properly already, so we're returning false to avoid calling remove_subscriber() again for them? It's a bit unclear to me how they can still be in the subscribers collection though, can you elaborate on that?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Indeed, I could actually .expect() here, since if the structure is consistent the subscribers should always be in self.sinks. The check here is superfluous.
I could refactor to:

let&(ref sink,ref filter) = self.sinks.get(&subscriber).expect("subscribers returned from self.listeners are always in self.sinks; qed");let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}

or

ifletSome(&(ref sink,ref filter)) = matchself.sinks.get(&subscriber){
let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}}

Which one do you prefer?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Oh, actually I can't since .get() borrows immutably, so remove_subscriber has to be outside of the scope.

assert_eq!(notifications.listeners.len(), 2);
assert_eq!(notifications.wildcard_listeners.len(), 1);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The channels are closed here, correct?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes and since the receiving end is dropped, sending to such channel will trigger an error.

sink
.sink_map_err(|e| warn!("Error sending notifications: {:?}", e))
.send_all(stream)
// we ignore the resulting Stream (if the first stream is over we are unsubscribed)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

…is over…? Do you mean …is closed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

is over as in is finished/is done. Which means that the stream will not emit any more items.

/// Drain committed changes to an iterator.
///
/// Panics:
/// Will panic if there are any uncommitted prospective changes.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is "prospective" sort of like "pending"?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, those are changes that can still be easily discarded. You can see example of usage inside block builder:

  1. We run a transaction
  2. It produces a set of prospective changes
  3. If we detect that it's somehow invalid we discard the prospective changes
  4. If we accept the transaction we commit prospective changes.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?
self.storage_notifications.lock()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are you sure that it should be called before transaction is committed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. Moved the notification after commit and also guarded by the same if as block import notification.

@gavofyorkgavofyork left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Aside from the minor comment

@@ -0,0 +1,267 @@
// Copyright 2017 Parity Technologies (UK) Ltd.
// This file is part of Polkadot.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Substrate, not Polkadot :)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed.

@gavofyorkgavofyork removed the A0-please_review Pull request needs code review. label Aug 1, 2018
@svyatonik
svyatonik merged commit 757a721 into masterAug 1, 2018
@svyatonik
svyatonik deleted the td-storage-events branch August 1, 2018 12:29
dvdplm added a commit that referenced this pull request Aug 1, 2018
* master:
Collator for the "adder" (formerly basic-add) parachain and various small fixes (#438)
Storage changes subscription (#464)
Wasm execution optimizations (#466)
Fix the --key generation (#475)
Fix typo in service.rs (#472)
Fix session phase in early-exit (#453)
Make ping unidirectional (#458)
Update README.adoc
gavofyork pushed a commit that referenced this pull request Aug 10, 2018
* Initial implementation of storage events.
* Attaching storage events.
* Expose storage modification stream over RPC.
* Use FNV for hashing small keys.
* Fix and add tests.
* Swap alias and RPC name.
* Fix demo.
* Addressing review grumbles.
* Fix comment.
liuchengxu pushed a commit to autonomys/substrate that referenced this pull request Jun 3, 2022
helin6 pushed a commit to boolnetwork/substrate that referenced this pull request Jul 25, 2023
* Bump release version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update dependency version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Modify changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Move changelog entries from added to changed
Sign up for freeto subscribe to this conversation on GitHub. Already have an account? Sign in.

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@tomusdrw@dvdplm@gavofyork@svyatonik@jacogr
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
This repository was archived by the owner on Nov 15, 2023. It is now read-only.

Storage changes subscription - #464

Merged
svyatonik merged 10 commits into
masterfrom
td-storage-events
Aug 1, 2018
Merged

Storage changes subscription#464
svyatonik merged 10 commits into
masterfrom
td-storage-events

Conversation

@tomusdrw

Copy link
Copy Markdown
Contributor
  • Exposing storage changes via RPC pub-sub
  • Allows one to subscribe to particular storage keys

CC @jacogr

@tomusdrwtomusdrw added A0-please_review Pull request needs code review. M6-rpcapi labels Jul 31, 2018
Comment threadsubstrate/rpc/src/chain/mod.rs Outdated
#[pubsub(name = "chain_newHead")] {
/// New head subscription
#[rpc(name = "subscribe_newHead")]
#[rpc(name = "subscribe_newHead", alias = ["chain_subscribeNewHead", ])]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Great. Maybe it is better to have the "old" as the alias and the "new" one as the name?

(e.g. chain_subscribeNewHead is actually probably the preferred one to match since it aligns with other RPCs, my gut tells me the preferred one should be the "default")

@dvdplmdvdplm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good stuff. I can't say I understand it all, but overall very readable and interesting.

Comment threadsubstrate/client/src/client.rs Outdated

/// Get storage changes event stream.
///
/// Passing `None` as keys subscribes to all possible keys

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

typo: …as keys should be …as key

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rephrased the whole sentence, hope it's clearer now.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some(storage_update) = storage_update {
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't we emit a re-org event with some meta data about what changed and let interested subscribers re-fetch?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, that's possible, although it's not easy to get storage changes for blocks that are already imported/executed. Re-fetching changes would mean to re-execute the blocks, but I suppose based on the filter_keys we could just return all the storage values in the re-orged blocks.

subscribers.extend(listeners.iter());
}

if has_wildcard || listeners.is_some() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't you check for !susbscribers.is_empty() here? Or is it faster to do it this way?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm actually interested only in subscribers for that particular key. So if there is a set of changes:

[(1, Some(2)), (2, None), (3, Some(4))]

but we have no wildcard_listeners and only listener for key=1, changes vector will only contain [(StorageKey(1), Some(StorageData(2))]

filter: filter.clone(),
})).is_err()
},
None => false,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If I read this right, if the subscriber is gone when we get here, then we assume they've been removed properly already, so we're returning false to avoid calling remove_subscriber() again for them? It's a bit unclear to me how they can still be in the subscribers collection though, can you elaborate on that?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Indeed, I could actually .expect() here, since if the structure is consistent the subscribers should always be in self.sinks. The check here is superfluous.
I could refactor to:

let&(ref sink,ref filter) = self.sinks.get(&subscriber).expect("subscribers returned from self.listeners are always in self.sinks; qed");let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}

or

ifletSome(&(ref sink,ref filter)) = matchself.sinks.get(&subscriber){
let result = sink.unbounded_send((hash.clone(),StorageChangeSet{changes: changes.clone(),filter: filter.clone(),}));if result.is_err(){self.remove_subscriber(subscriber);}}

Which one do you prefer?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Oh, actually I can't since .get() borrows immutably, so remove_subscriber has to be outside of the scope.

assert_eq!(notifications.listeners.len(), 2);
assert_eq!(notifications.wildcard_listeners.len(), 1);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The channels are closed here, correct?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes and since the receiving end is dropped, sending to such channel will trigger an error.

sink
.sink_map_err(|e| warn!("Error sending notifications: {:?}", e))
.send_all(stream)
// we ignore the resulting Stream (if the first stream is over we are unsubscribed)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

…is over…? Do you mean …is closed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

is over as in is finished/is done. Which means that the stream will not emit any more items.

/// Drain committed changes to an iterator.
///
/// Panics:
/// Will panic if there are any uncommitted prospective changes.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is "prospective" sort of like "pending"?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, those are changes that can still be easily discarded. You can see example of usage inside block builder:

  1. We run a transaction
  2. It produces a set of prospective changes
  3. If we detect that it's somehow invalid we discard the prospective changes
  4. If we accept the transaction we commit prospective changes.

Comment threadsubstrate/client/src/client.rs Outdated
if let Some((storage_update, changes)) = storage_update {
transaction.update_storage(storage_update)?;
// TODO [ToDr] How to handle re-orgs? Should we re-emit all storage changes?
self.storage_notifications.lock()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are you sure that it should be called before transaction is committed?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. Moved the notification after commit and also guarded by the same if as block import notification.

@gavofyorkgavofyork left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Aside from the minor comment

@@ -0,0 +1,267 @@
// Copyright 2017 Parity Technologies (UK) Ltd.
// This file is part of Polkadot.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Substrate, not Polkadot :)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed.

@gavofyorkgavofyork removed the A0-please_review Pull request needs code review. label Aug 1, 2018
@svyatonik
svyatonik merged commit 757a721 into masterAug 1, 2018
@svyatonik
svyatonik deleted the td-storage-events branch August 1, 2018 12:29
dvdplm added a commit that referenced this pull request Aug 1, 2018
* master:
Collator for the "adder" (formerly basic-add) parachain and various small fixes (#438)
Storage changes subscription (#464)
Wasm execution optimizations (#466)
Fix the --key generation (#475)
Fix typo in service.rs (#472)
Fix session phase in early-exit (#453)
Make ping unidirectional (#458)
Update README.adoc
gavofyork pushed a commit that referenced this pull request Aug 10, 2018
* Initial implementation of storage events.
* Attaching storage events.
* Expose storage modification stream over RPC.
* Use FNV for hashing small keys.
* Fix and add tests.
* Swap alias and RPC name.
* Fix demo.
* Addressing review grumbles.
* Fix comment.
liuchengxu pushed a commit to autonomys/substrate that referenced this pull request Jun 3, 2022
helin6 pushed a commit to boolnetwork/substrate that referenced this pull request Jul 25, 2023
* Bump release version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Update dependency version to v0.18.0
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Modify changelog
Signed-off-by: Alexandru Vasile <alexandru.vasile@parity.io>
* Move changelog entries from added to changed
Sign up for freeto subscribe to this conversation on GitHub. Already have an account? Sign in.

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@tomusdrw@dvdplm@gavofyork@svyatonik@jacogr