You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Following the storage node design, the distributor node CLI is currently a separate package, independent from joystream-cli, but sharing some similarities in the code (both are build with OCLIF framework, both manage api and accounts kind of similarly etc.).
The CLI is also the root of the project, it is used both to run the actual node and to perform some on-chain operations using the substrate api.
Currently the documentation of all supported commands can be found HERE
Commands were designed with scripability in mind and can easily be tested using bash scripts (see current example script)
Further developments:
There are a few extrinsics commands missing, like set-family-metadata or remove-bucket-operator, which weren't initially supported by the runtime (see: Storage v2 distribution - progress tracking #2543)
It's currently unclear to what extend the state inspection commands should be implemented, since the assumption is that operators and lead will most likely have to get familiar with the graphql playground at some point anyway in order to be able to inspect the data they're interested in (see: Storage v2 distribution - progress tracking #2543). On the other hand it should be relatively easy to integrate some basic query-node and chain queries to the CLI and it would reduce the number of tools required to perform given role duties and make the experience much more convinient.
Idea: Merging with joystream-cli
It's certainly possible to merge the current distributor CLI commands to joystream-cli and perhaps also improve the latter in the process. It should be considered whether it would make more sense to have "one CLI for everyone" or multiple, role-specific CLI's.
Idea: Abstracting out common CLI functionality
Currently, with 3 different CLIs there is some redundancy/inconsistencies in the code w.r.t. how keys, api connections, validation etc. are handled (this issue goes even further if we consider integration tests and brings back the old idea of joystream-js, which never materialized due to the fact that its purpose wasn't yet clear at the time, infrastructure was changing rapidly and the process of introducing such library in the monorepo environment turned out to be non-trivial and require structural changes affecting the entire repository).
Configuration
The distributor node CLI and the node itself are using a single YAML/JSON configuration file, the path to which:
can be provided through --config or -c flag supported by all the commands
can be provided through CONFIG_PATH env variable
is by default assumed to be config.yml file in the current working directory
The current config file has a following structure:
Most of the properties are self-explainatory, the ones that may not be immediately obvious include:
log.file - minimum log level of logs that end up in the logfile (the logfile location is determined by {directories.logs}/logs.json). Standard npm logging levels are supported (https://www.npmjs.com/package/winston#logging-levels)
log.console - minimum log level of logs that are printed to the console (optional)
keys - array of keys that will be added to the CLI/node keyring. Currently provided in form of substrate uris, this will be changed to support other ways of providing keys (ie. path to a file)
buckets - buckets that will be served by the distributor node instance. This will be adjusted to support an 'all' configuration, allowing to track and serve all buckets currently assigned to the node operator (via the runtime)
storageLimit - the limit of how much storage can be utilized for data object caching purposes. G stands for gigabytes. (for details see: Content and caching policy)
Distributor node API
The API is described by an OpenAPI schema (see the current schema here) and implemented with ExpressJS.
It currently exposes one endpoint - /asset/{objectId}, which serves the file by a specified objectId.
There are multiple scenarios of how a distributor node may act upon a request:
Scenario 1: The requested object is available in the distributor node filesystem (cache):
In this case:
The object TLRU cache position is updated (see Content and caching policy for more details)
The object is served via send library which supports partial responses (Ranges), conditional-GET negotiation (If-Match, If-Unmodified-Since, If-None-Match, If-Modified-Since) etc. through HTTP headers
TODO: There should be an additional check whether the object is still supposed to be served by the distributor node. If not, it should be dropped from the cache either before or after serving the request (since checking this would require communication with the query-node, it may be more suitable to do this after the request is served, to avoid any unnecessary delays)
TODO: There should be an additional HTTP header sent, indicating that the distributor node cache is beeing used (ie. X-Cache: HIT)
Scenario 2: The object is not yet cached, but is currently beeing fetched from the storage node
In this case the request is forwarded to the storage node that the object is beeing downloaded from.
The reason why the request is not processed by streaming the data from file while it is still beeing fetched (which was my initial approach) is that the flow of how, for example, a video is beeing fetched by the browser, usually consists of making multiple partial (Range) requests, often targetting the very end of the file (most likely to read some of its metadata first). These kind of requests will be acted upon faster when forwarded to a storage node which already has the full data available.
TODO: There should be an additional HTTP header sent, indicating that the distributor node cache was missed, or more specifically, that the object is still being fetched (ie. X-Cache: MISS / X-Cache: FETCHING)
Scenario 3: The object is neither cached not currently beeing fetched
In this case the distributor node is making a request to query node to fetch details of the requested object, which include: content hash, object size, who is supposed to store the object, who is supposed to distribute it etc. It then proceeds to one of the following scenarios:
Scenario 3.1: The object doesn't exist
Node responds with HTTP 404 and a message
Scenario 3.2: The object is not distributed by the node
Node responds with HTTP 400 and a message
Scenario 3.3: The request was valid, the node needs to fetch the object
This triggers a pretty complex process of fetching the data object from storage node, which is described in detail in the Fetching data from storage node section below.
Once the storage node from which the object will be fetched is found, the request is handled in a way analogous to the one described in Scenario 2.
TODO: There should be an additional HTTP header sent, indicating that the distributor node cache was missed (ie. X-Cache: MISS)
Further plans
support HEAD requests on data object endpoint or a separate endpoint to inspect object cache state without forcing the node to fetch / serve the actual obejct
add a status / healthcheck endpoint, returning some basic information about node's current status (for example, used / free storage, number of objects stored / beeing fetched etc.)
add an endpoint informing about supported buckets
Idea: authorized api endpoints to manage some of the node configuration (served buckets, storage limit), force a filesystem check, retry a data object download or triggering other actions that could be then performed asynchronously while the node is running (alternatively, some of those may be part of the CLI and only runnable on the actual machine that runs the distributor node itself)
possibly expose other, potentially useful node state (TBD)
Fetching data from storage nodes
Finding nearby storage nodes:
In order to limit the number of requests beeing made on cache miss and the time it takes to fetch a new object in this scenario, the distributor node needs to keep some state about which storage endpoints will probably respond to a request most quickly.
This can be partially solved by making use of the metadata provided by storage node operators, which may include some approximate geographic coordinates (see Metadata standards section), which may provide some estimation of which nodes will likely respond faster (ie. a node that is supposedly 100 kilometers away will most likely respond faster that the one 10000 kilometers away). However, because the approach of using metadata geographic coordinates is very limited and it's possible that most storage providers will choose not to expose such information, the distributor node instead uses a different approach.
Currently the distributor node periodically (every x seconds) makes requests to all active storage provider endpoints (fetched from the query node) and measures their average response times. This is done asynchronously and independently of any incoming requests. By using the mechanism of queueing the requests with a relatively small concurrency limit, this process is relatively low-cost and provides a pretty good estimation on which nodes will likely be the best candidates for fetching data objects during a cache miss.
Data object fetching flow
A cache miss, as described in Scenario 3.3 in API section, triggers the folowing flow:
First, the endpoints of storage providers that are supposed to store the given object are ordered by the mean response time (the process of obtaining it is described in the previous section)
The requests are then sent to the storage endpoints, starting from the ones with lowest mean response time. Those are data object availibitly check requests, which are meant only to determine whether a given storage node indeed has the data object available. They are queued using a choosen concurrency (10 at the current time).
TODO: Currently there is no explicit way of checking whether the storage node has a given object available. The distributor node achieves this by sending a GET request with Range: 'bytes=0-0' header, effectively requesting just the first byte of the data to limit the potential waste of bandwidth on the storage node side. Ideally, there would be a HEAD request supported by the storage node for this purpose.
As soon as any storage provider endpoint confirms the availability of the object, the availabilityCheckQueue is temporarly stopped and a request is made to fetch the full data from the selected provider. If multiple storage providers confirm the availability of the object at roughly the same time, the fetch / download requests will be added to a queue (which uses a concurrency of 1), allowing the distributor node to instantly try a different provider in case the actual data fetch request to a given provider fails. The process continues until a storage node that succesfully responds to this request is found.
Once the storage node succesfully responds with the data, all other requests w.r.t. that content are stopped and the node begins to write the data into its filesystem. Any errors at this point (unexpected data size, stream errors) will mean a failure to fetch the data object, causing the content to not be stored and the whole process of fetching the object to potentially be repeated later.
TODO: Currently no retry attempts are actually happening during such failure. The object is dropped entirely and the node will only try to fetch it again if there is another request for it. This seems to be an issue and could also potentially be expolited by a malicious/faulty storage node, so some adjustments may be needed (possibly including keeping some state about faulty storage nodes in order to quickly bypass them).
Fetching data from other distributor nodes (?)
It is not clear whether distributor-to-distributor communication should also be utilized for the purpose of fetching missing data objects. It certainly adds a lot of complexity, which we probably want to avoid in the first iteration of Storage v2.
The state
Most of the state, including both the part that's only stored in node's memory and doesn't need to be persisted across restarts and the part that is persisted in the storage (currently via an asynchronously updated json file), is handled via an "intermediary" StateCacheService. This is to faciliate the potential migration to other state management approaches, like using a Redis database etc.
The current node state includes:
Memory
pendingDownloadsByContentHash map - mapps content hashes to information about pending downloads (data object fetching attempts). Initially it stored the information in form of simple json object, currently it just stores the reference to downloadPromise which resolves to a response from a storage node that has been chosen to be the data source. The response object (AxiosResponse) contains pretty much all of the useful information, like the storage node endpoint, the actual data stream, response headers etc.
contentHashByObjectId - a simple map of dataObjectId -> contentHash
storageNodeEndpointDataByEndpoint - currently stores average mean response times mapped by storage nodes endpoints (see: Finding nearby storage nodes)
Memory + persistent storage
contentLruDataByHash - stores TLRU-cache-related information by content hash. This includes object TTU (as a function of data size) and current validUntil value (For details see Content and caching policy section)
mimeTypeByContentHash - a map that stores the relationship between content hash and the data mimeType (as determined by the distributor node)
TODO: The exact shape of the function of how object size should affect TTU is not exactly clear yet.
Storage limit
The node operator can configure the storage limit - the amount of storage that will be devoted to the cache. The node will always try to maximally utilize the storage up to the provided limit, which means no data objects are currently removed just because they are stale or no longer supported. The data objects are only removed to make place for the new data objects, in case there is currently no free space available (accoriding to the limit). This also means it is not possible for srored objects size to go beyond this limit.
Further developments:
It would make sense to make some adjustments, so that the stale data objects (past their lifetime) and objects that are no longer supposed to be served by the node are dropped as soon as possible. This would make the cache more flexible and reduce the average amount of storage beeing used (in cases the limit is relatively high). The amount of storage used in that case would also become a good indicator of how much storage the distributor node really "needs" depending on user activity. A hybrid approach where a storage limit is possible to be set, but not necessary for the cache to function properly should be implemented.
Startup, cleanup, data integrity
Startup and cleanup functions are an important piece of the distributor node design.
During a startup, some data integrity checks are performed, which potentially may later be supported also while the node is still running. An additional, more thorough check, including calculataion of data object hashes, re-reading files mimeType metadata etc. may also be added in the future.
Current startup / cleanup logic includes:
Startup:
Load state stored by StateCacheService from the filesystem
Fetch data about all objects that the node is supposed to distribute based on the configured buckets (or, in the near future, buckets assigned to a given operator) from the query node
Iterate over objects in the filesystem and:
drop no longer supported objects
drop objects of invalid size (TODO: This will change, as the node will actually try to resume fetching the data objects that haven't been fully fetched yet)
determine the mimeType of objects for which the mimeTypeByContentHash entry is missing
recreate contentHashByObjectId map
TODO: The node should also resume any pending data object fetches / downloads, by trying to fetch objects that are either partially fetched or not fetched at all, but are part of the TLRU cache data. This is currently not supported.
Cleanup:
There are multiple ways the node can be shut down, obviously not all of them allow performing any cleanup at all (ie. power outage, forced process kill etc.). Currently the node is trying to make the most graceful exit possible given the circumstances.
The node will always try to make sure all the stored data, like logs and StateCacheService data, is flushed to the disk before exiting.
Idea: If the node process is terminated in a "non-instant" way (for example, through CTRL+C), it may also try to wait until all the current requests are served and all the data objects beeing fetched from the storage nodes are fully fetched and stored in the local fs. This shouldn't take more than 30 second though, because after that time the process may be forceully killed.
Logging
The distributor node, just like the storage node, is using winston library for logging.
I order to test the Elastic Stack integration, I created a docker-compose file for running the stack locally along with the distributor node itself. The flow of the logs in that case is following (which is just one of many possible paths): Distributor node (winston) -> logfile -> FileBeat -> ElasticSearch -> Kibana
Metadata standards
The current metadata standards are described by following a set of protobuf messages:
message GeoCoordiantes {
required float latitude = 3;
required float longitude = 4;
}
message NodeLocationMetadata {
optional string country_code = 1; // ISO 3166-1 alpha-2 country code (2 letters)
optional string city = 2; // City name
optional GeoCoordiantes coordinates = 3; // Geographic coordinates
}
message StorageBucketOperatorMetadata {
optional string endpoint = 1; // Root storage node endpoint (ie. https://example.com/storage)
optional NodeLocationMetadata location = 2; // Information about node's phisical location
optional string extra = 3; // Additional information about the node / node operator
}
message DistributionBucketOperatorMetadata {
optional string endpoint = 1; // Root distribution node endpoint (ie. https://example.com/distribution)
optional NodeLocationMetadata location = 2; // Information about node's phisical location
optional string extra = 3; // Additional information about the node / node operator
}
message DistributionBucketFamilyMetadata {
optional string region = 1; // ID / name of the region covered by the distribution family (ie. us-east-1). Should be unique.
optional string description = 2; // Additional, more specific description of the region
repeated GeoCoordiantes boundary = 3; // Geographical boundary of the region, defined as polygon through array of coordinates
}
In practice, only a few of those metadata fields may end up actually beeing used, but my reasoning was that allowing the operators to provide some optional information about the node may have some utility both in terms of debugging and optimizing some processes related to grouping sorting nodes etc., for example, based on location.
The most critical part of the metadata are node endpoints - in that regard, the metadata is the only viable source of infromation about how to access given a node.
Distribution bucket families
It is important that the frontend applications (ie. Atlas) can easily identify which distribution bucket families should be preffered for fetching the assets, based on user's location and connection.
There are multiple ways of how this can be achieved and the current metadata standard is quite flexible in that regard, not enforcing any specific semantics on, for example, what is considered a region.
There are, however, a few ideas I had in mind while designing the standard:
Using client's geolocation and geographic regions:
the frontend app should be able to determine most suitable regions for the client by using, for example, browser's Geolocation API and and some standard of how to translate this location to a family bucket region
in order to make it possible to include the information about some agreed-upon standard of how to transalte client's geolocation to most suitable family bucket regions, I included an optional boundary field in the metadata standard which can be used to provide the exact coordinates of a polygon that describes a geographic region covered by a given family
Using latency test:
The forentend application can perform a simple latency test, pinging some endpoints that are known to be associated with a given region (for example, see: https://www.cloudping.info/) in order to determine which regions will be the best fit for a client.
I was considering whether the family bucket metadata should include a field like pingTargets that would allow specifying those endpoints, but decided to drop it to reduce the complexity. After some reconsideration I think it may be actually be a good idea to include it.
Query node
Mappings
There is currently a pending PR with the query-node storage v2 schema and mappings: #2515
Missing parts:
Setting distribution bucket family metadata (wasn't supported by the runtime yet)
Removing distribution bucket operator (wasn't supported by the runtime yet)
Depending on how Storage Node syncing will be handled, there may be a need to adjust the schema and mappings to faciliate fetching only the delta of the storage v2 state based on provided timestamp. One way to do this may be adding a date field to a StorageBag entity that indicate when was the last time the data objects inside a given bag changed. This could greatly reduce the amount of objects beeing fetched from the query node per each storage node state update period.
Unresolved issues:
There seems to currently be an issue with multiple many-to-many relationships (ie. bag->distributedBy, bag->storedBy, storageBucket->storedBag, distributionBucket->distributedBags) that causes only one side of the relationship to be work when querying the data from the node (for example: bag->storedBy returns the correct bucket, but storageBucket->storedBags returns empty set. For distribution buckets it's the other way around). Potentially related issue: Problem when querying ManyToMany relationships in hydra-3.1.0-alpha.7 hydra#448. This can be temporarly sidestepped by introcuding an intermediary entity like DistributionBucketToBagRelationship with 2 many-to-one relationships: one to a bucket and one to a bag.
I wasn't "manually" setting createdAt and updatedAt dates with the hope that this wouldn't be needed in the near future, as it has already been solved in the newest version of Hydra. Unfortunately the update to the newest version is currently blocked by Problem when querying ManyToMany relationships in hydra-3.1.0-alpha.7 hydra#448, so those dates may be temporarly unreliable.
Integration with distributor node
For maximum efficiency, the processor part of the query node (ideally containing only the storage v2 related mappings) should probably be run either on the same machine as the distributor node itself, or on a machine that's almost equally quickly accessible.
The access to the data can be even further optimized if the distributor node connects to the processor database directly, instead of using the graphql-server interface. This needs to be further evaluated in terms of pros and cons.
CLI
Following the storage node design, the distributor node CLI is currently a separate package, independent from
joystream-cli, but sharing some similarities in the code (both are build with OCLIF framework, both manage api and accounts kind of similarly etc.).The CLI is also the root of the project, it is used both to run the actual node and to perform some on-chain operations using the substrate api.
Currently the documentation of all supported commands can be found HERE
Commands were designed with scripability in mind and can easily be tested using bash scripts (see current example script)
Further developments:
set-family-metadataorremove-bucket-operator, which weren't initially supported by the runtime (see: Storage v2 distribution - progress tracking #2543)joystream-cliIt's certainly possible to merge the current distributor CLI commands to
joystream-cliand perhaps also improve the latter in the process. It should be considered whether it would make more sense to have "one CLI for everyone" or multiple, role-specific CLI's.Currently, with 3 different CLIs there is some redundancy/inconsistencies in the code w.r.t. how keys, api connections, validation etc. are handled (this issue goes even further if we consider integration tests and brings back the old idea of
joystream-js, which never materialized due to the fact that its purpose wasn't yet clear at the time, infrastructure was changing rapidly and the process of introducing such library in the monorepo environment turned out to be non-trivial and require structural changes affecting the entire repository).Configuration
The distributor node CLI and the node itself are using a single YAML/JSON configuration file, the path to which:
--configor-cflag supported by all the commandsCONFIG_PATHenv variableconfig.ymlfile in the current working directoryThe current config file has a following structure:
Most of the properties are self-explainatory, the ones that may not be immediately obvious include:
log.file- minimum log level of logs that end up in the logfile (the logfile location is determined by{directories.logs}/logs.json). Standard npm logging levels are supported (https://www.npmjs.com/package/winston#logging-levels)log.console- minimum log level of logs that are printed to the console (optional)keys- array of keys that will be added to the CLI/node keyring. Currently provided in form of substrate uris, this will be changed to support other ways of providing keys (ie. path to a file)buckets- buckets that will be served by the distributor node instance. This will be adjusted to support an'all'configuration, allowing to track and serve all buckets currently assigned to the node operator (via the runtime)storageLimit- the limit of how much storage can be utilized for data object caching purposes.Gstands for gigabytes. (for details see: Content and caching policy)Distributor node API
The API is described by an OpenAPI schema (see the current schema here) and implemented with ExpressJS.
It currently exposes one endpoint -
/asset/{objectId}, which serves the file by a specifiedobjectId.There are multiple scenarios of how a distributor node may act upon a request:
Scenario 1: The requested object is available in the distributor node filesystem (cache):
In this case:
sendlibrary which supports partial responses (Ranges), conditional-GET negotiation (If-Match,If-Unmodified-Since,If-None-Match,If-Modified-Since) etc. through HTTP headersTODO: There should be an additional check whether the object is still supposed to be served by the distributor node. If not, it should be dropped from the cache either before or after serving the request (since checking this would require communication with the query-node, it may be more suitable to do this after the request is served, to avoid any unnecessary delays)
TODO: There should be an additional HTTP header sent, indicating that the distributor node cache is beeing used (ie.
X-Cache: HIT)Scenario 2: The object is not yet cached, but is currently beeing fetched from the storage node
In this case the request is forwarded to the storage node that the object is beeing downloaded from.
The reason why the request is not processed by streaming the data from file while it is still beeing fetched (which was my initial approach) is that the flow of how, for example, a video is beeing fetched by the browser, usually consists of making multiple partial (
Range) requests, often targetting the very end of the file (most likely to read some of its metadata first). These kind of requests will be acted upon faster when forwarded to a storage node which already has the full data available.TODO: There should be an additional HTTP header sent, indicating that the distributor node cache was missed, or more specifically, that the object is still being fetched (ie.
X-Cache: MISS/X-Cache: FETCHING)Scenario 3: The object is neither cached not currently beeing fetched
In this case the distributor node is making a request to query node to fetch details of the requested object, which include: content hash, object size, who is supposed to store the object, who is supposed to distribute it etc. It then proceeds to one of the following scenarios:
Scenario 3.1: The object doesn't exist
Node responds with
HTTP 404and a messageScenario 3.2: The object is not distributed by the node
Node responds with
HTTP 400and a messageScenario 3.3: The request was valid, the node needs to fetch the object
This triggers a pretty complex process of fetching the data object from storage node, which is described in detail in the
Fetching data from storage nodesection below.Once the storage node from which the object will be fetched is found, the request is handled in a way analogous to the one described in
Scenario 2.TODO: There should be an additional HTTP header sent, indicating that the distributor node cache was missed (ie.
X-Cache: MISS)Further plans
HEADrequests on data object endpoint or a separate endpoint to inspect object cache state without forcing the node to fetch / serve the actual obejctstatus/healthcheckendpoint, returning some basic information about node's current status (for example, used / free storage, number of objects stored / beeing fetched etc.)Fetching data from storage nodes
Finding nearby storage nodes:
In order to limit the number of requests beeing made on cache miss and the time it takes to fetch a new object in this scenario, the distributor node needs to keep some state about which storage endpoints will probably respond to a request most quickly.
This can be partially solved by making use of the metadata provided by storage node operators, which may include some approximate geographic coordinates (see Metadata standards section), which may provide some estimation of which nodes will likely respond faster (ie. a node that is supposedly 100 kilometers away will most likely respond faster that the one 10000 kilometers away). However, because the approach of using metadata geographic coordinates is very limited and it's possible that most storage providers will choose not to expose such information, the distributor node instead uses a different approach.
Currently the distributor node periodically (every
xseconds) makes requests to all active storage provider endpoints (fetched from the query node) and measures their average response times. This is done asynchronously and independently of any incoming requests. By using the mechanism of queueing the requests with a relatively small concurrency limit, this process is relatively low-cost and provides a pretty good estimation on which nodes will likely be the best candidates for fetching data objects during a cache miss.Data object fetching flow
A cache miss, as described in
Scenario 3.3inAPIsection, triggers the folowing flow:First, the endpoints of storage providers that are supposed to store the given object are ordered by the mean response time (the process of obtaining it is described in the previous section)
The requests are then sent to the storage endpoints, starting from the ones with lowest mean response time. Those are data object availibitly check requests, which are meant only to determine whether a given storage node indeed has the data object available. They are queued using a choosen concurrency (
10at the current time).TODO: Currently there is no explicit way of checking whether the storage node has a given object available. The distributor node achieves this by sending a
GETrequest withRange: 'bytes=0-0'header, effectively requesting just the first byte of the data to limit the potential waste of bandwidth on the storage node side. Ideally, there would be aHEADrequest supported by the storage node for this purpose.As soon as any storage provider endpoint confirms the availability of the object, the
availabilityCheckQueueis temporarly stopped and a request is made to fetch the full data from the selected provider. If multiple storage providers confirm the availability of the object at roughly the same time, thefetch/downloadrequests will be added to a queue (which uses a concurrency of1), allowing the distributor node to instantly try a different provider in case the actual data fetch request to a given provider fails. The process continues until a storage node that succesfully responds to this request is found.Once the storage node succesfully responds with the data, all other requests w.r.t. that content are stopped and the node begins to write the data into its filesystem. Any errors at this point (unexpected data size, stream errors) will mean a failure to fetch the data object, causing the content to not be stored and the whole process of fetching the object to potentially be repeated later.
TODO: Currently no retry attempts are actually happening during such failure. The object is dropped entirely and the node will only try to fetch it again if there is another request for it. This seems to be an issue and could also potentially be expolited by a malicious/faulty storage node, so some adjustments may be needed (possibly including keeping some state about faulty storage nodes in order to quickly bypass them).
Fetching data from other distributor nodes (?)
It is not clear whether distributor-to-distributor communication should also be utilized for the purpose of fetching missing data objects. It certainly adds a lot of complexity, which we probably want to avoid in the first iteration of Storage v2.
The state
Most of the state, including both the part that's only stored in node's memory and doesn't need to be persisted across restarts and the part that is persisted in the storage (currently via an asynchronously updated
jsonfile), is handled via an "intermediary"StateCacheService. This is to faciliate the potential migration to other state management approaches, like using aRedisdatabase etc.The current node state includes:
Memory
pendingDownloadsByContentHashmap - mapps content hashes to information about pending downloads (data object fetching attempts). Initially it stored the information in form of simple json object, currently it just stores the reference todownloadPromisewhich resolves to a response from a storage node that has been chosen to be the data source. The response object (AxiosResponse) contains pretty much all of the useful information, like the storage node endpoint, the actual data stream, response headers etc.contentHashByObjectId- a simple map ofdataObjectId->contentHashstorageNodeEndpointDataByEndpoint- currently stores average mean response times mapped by storage nodes endpoints (see: Finding nearby storage nodes)Memory + persistent storage
contentLruDataByHash- stores TLRU-cache-related information by content hash. This includes object TTU (as a function of data size) and currentvalidUntilvalue (For details see Content and caching policy section)mimeTypeByContentHash- a map that stores the relationship between content hash and the datamimeType(as determined by the distributor node)Content and caching policy
The current caching policy is a form of
TLRU(https://en.wikipedia.org/wiki/Cache_replacement_policies#Time_aware_least_recently_used_(TLRU)) where theTTUcomponent is a function of data object size, such that the small data objects have higherTTU, while big data objects have lowerTTU, making the smaller data objects more favorable to store over longer periods of time.TODO: The exact shape of the function of how object size should affect
TTUis not exactly clear yet.Storage limit
The node operator can configure the storage limit - the amount of storage that will be devoted to the cache. The node will always try to maximally utilize the storage up to the provided limit, which means no data objects are currently removed just because they are stale or no longer supported. The data objects are only removed to make place for the new data objects, in case there is currently no free space available (accoriding to the limit). This also means it is not possible for srored objects size to go beyond this limit.
Further developments:
It would make sense to make some adjustments, so that the stale data objects (past their lifetime) and objects that are no longer supposed to be served by the node are dropped as soon as possible. This would make the cache more flexible and reduce the average amount of storage beeing used (in cases the limit is relatively high). The amount of storage used in that case would also become a good indicator of how much storage the distributor node really "needs" depending on user activity. A hybrid approach where a storage limit is possible to be set, but not necessary for the cache to function properly should be implemented.
Startup, cleanup, data integrity
Startup and cleanup functions are an important piece of the distributor node design.
During a startup, some data integrity checks are performed, which potentially may later be supported also while the node is still running. An additional, more thorough check, including calculataion of data object hashes, re-reading files
mimeTypemetadata etc. may also be added in the future.Current startup / cleanup logic includes:
Startup:
StateCacheServicefrom the filesystemmimeTypeof objects for which themimeTypeByContentHashentry is missingcontentHashByObjectIdmapTLRUcache data. This is currently not supported.Cleanup:
There are multiple ways the node can be shut down, obviously not all of them allow performing any cleanup at all (ie. power outage, forced process kill etc.). Currently the node is trying to make the most graceful exit possible given the circumstances.
The node will always try to make sure all the stored data, like logs and
StateCacheServicedata, is flushed to the disk before exiting.Idea: If the node process is terminated in a "non-instant" way (for example, through
CTRL+C), it may also try to wait until all the current requests are served and all the data objects beeing fetched from the storage nodes are fully fetched and stored in the local fs. This shouldn't take more than 30 second though, because after that time the process may be forceully killed.Logging
The distributor node, just like the storage node, is using
winstonlibrary for logging.This is a very popular NodeJS logging library, it's also easy to integrate it with the Elastic stack (https://www.elastic.co/guide/en/ecs-logging/nodejs/current/winston.html) and an ExpressJS server. It allows specifying multiple targets for the logs to be directed to (currenty the distributor node logs are stored in a file and outputted to a console), supports multiple log levels (see: https://www.npmjs.com/package/winston#logging-levels), allows managing multiple log formats etc.
I order to test the Elastic Stack integration, I created a docker-compose file for running the stack locally along with the distributor node itself. The flow of the logs in that case is following (which is just one of many possible paths):
Distributor node (winston)->logfile->FileBeat->ElasticSearch->KibanaMetadata standards
The current metadata standards are described by following a set of protobuf messages:
In practice, only a few of those metadata fields may end up actually beeing used, but my reasoning was that allowing the operators to provide some optional information about the node may have some utility both in terms of debugging and optimizing some processes related to grouping sorting nodes etc., for example, based on location.
The most critical part of the metadata are node endpoints - in that regard, the metadata is the only viable source of infromation about how to access given a node.
Distribution bucket families
It is important that the frontend applications (ie. Atlas) can easily identify which distribution bucket families should be preffered for fetching the assets, based on user's location and connection.
There are multiple ways of how this can be achieved and the current metadata standard is quite flexible in that regard, not enforcing any specific semantics on, for example, what is considered a
region.There are, however, a few ideas I had in mind while designing the standard:
Using client's geolocation and geographic regions:
regioncan mean a relatively large geographic area. Perhaps semantic similar to that ofAWSregions may be used (https://docs.aws.amazon.com/AWSEC2/latest/UserGuide/using-regions-availability-zones.html#concepts-available-regions)boundaryfield in the metadata standard which can be used to provide the exact coordinates of a polygon that describes a geographic region covered by a given familyUsing latency test:
The forentend application can perform a simple latency test, pinging some endpoints that are known to be associated with a given region (for example, see: https://www.cloudping.info/) in order to determine which regions will be the best fit for a client.
I was considering whether the family bucket metadata should include a field like
pingTargetsthat would allow specifying those endpoints, but decided to drop it to reduce the complexity. After some reconsideration I think it may be actually be a good idea to include it.Query node
Mappings
There is currently a pending PR with the query-node storage v2 schema and mappings: #2515
Missing parts:
StorageBagentity that indicate when was the last time the data objects inside a given bag changed. This could greatly reduce the amount of objects beeing fetched from the query node per each storage node state update period.Unresolved issues:
bag->distributedBy,bag->storedBy,storageBucket->storedBag,distributionBucket->distributedBags) that causes only one side of the relationship to be work when querying the data from the node (for example:bag->storedByreturns the correct bucket, butstorageBucket->storedBagsreturns empty set. For distribution buckets it's the other way around). Potentially related issue: Problem when querying ManyToMany relationships in hydra-3.1.0-alpha.7 hydra#448. This can be temporarly sidestepped by introcuding an intermediary entity likeDistributionBucketToBagRelationshipwith 2 many-to-one relationships: one to a bucket and one to a bag.createdAtandupdatedAtdates with the hope that this wouldn't be needed in the near future, as it has already been solved in the newest version of Hydra. Unfortunately the update to the newest version is currently blocked by Problem when querying ManyToMany relationships in hydra-3.1.0-alpha.7 hydra#448, so those dates may be temporarly unreliable.Integration with distributor node
For maximum efficiency, the processor part of the query node (ideally containing only the storage v2 related mappings) should probably be run either on the same machine as the distributor node itself, or on a machine that's almost equally quickly accessible.
The access to the data can be even further optimized if the distributor node connects to the processor database directly, instead of using the graphql-server interface. This needs to be further evaluated in terms of pros and cons.