Repository navigation
Generic Non-Blocking Task Management ("Queue") for discovery and nodes domains #329
Description
Activity
- addeddevelopmentStandard developmentStandard developmentdesignRequires designRequires designenhancementNew feature or requestNew feature or request
on Feb 9, 2022 There's a go implementation of persistent priority queue backed by leveldb here: https://github.com/beeker1121/goque. It can be used as a reference for this. Also a JS implementation here: https://github.com/eugeneware/level-q (the goque is probably more comprehensive).
Our priority queue needs to by default maintain order, because we do want to know the sorted list of jobs. But also allow us to add a special priority number on top.
From my imagination:
This reminds me of the indexing problem, where you can ask for a list of rows sorted by several columns. The first column would dictate the base sort, then subsequent columns would sort any ambiguous sub-orders.
Imagine we had 2 indexes. The first being your priority index using an arbitrary number, the second being the monotonic time index using
IdSortable. You could sort on the priority index first, then sort onIdSortablesecond.Maybe this then has a relationship to #188.
However this would only be for if you are looking up items. If we are streaming data from the level db that may be more complicated.
This can be done with a compound index. Prefix can be the priority number (lexinted), suffix can be IdSortable.
This means you can then stream results that are always ordered in terms of priority first then by time second.
Priority can start at 0 by default, and one can increment priorities depending on the origin of the task. Like tasks emitted by user wanting to lookup something can be set to a higher priority number.
We could do this directly by changing the queue domain key. But I'd suggest first solving the indexing issue in general first then building a compound index on top.
We discovered that the priority queue can also benefit from a uniqueness index creating a uniqueness constraint: #311 (comment)
This means that duplicate tasks cannot go into the priority queue. Not entirely sure if this is required because a queue can still say they should process the same task over and over.
We should have a concurrency bound in the queue. This means how many tasks should be executed at the same time. By default unbounded meaning all tasks gets executed immediately without waiting to be done.
For IO bound tasks, you might as well have unbounded concurrency. For CPU-bound it can be sent to the web worker pool which is bounded by core count. Battery usage optimisation may also affect our limit too.
A generic Queue class has been implemented here: 91287ab
This queue is not persistent, or a priority queue, however, it is designed to be a generic queue that can eventually be used in all places that require this functionality (including the Discovery Queue). The generic Queue can be refactored to meet this issue and #328 at some point in the future, potentially incorporating the DB.
- changed the title
[-]Refactor the Discovery Queue to be a Priority Queue prioritising new vertices[/-][+]Generic Non-Blocking Task Management ("Queue") for discovery and nodes domains[/+]on Apr 26, 2022 Renamed this issue to the general idea of non-blocking task management. It now has to solve for discovery, nodes management in terms of setting nodes, pinging nodes and garbage collection, as well as in relation to:
- Refactor error handling of failed Node Connections created from the Discovery domain #354
- Write tests to cover all possible states of the Discovery, NodeConnection, and discovered Node during the discovery process #349
- Supporting the "refresh" operation for
NodeGraphbuckets #345 - Discovery - revisiting Gestalt Vertices and error handling #328
- Asynchronous Promise Cancellation with Cancellable Promises, AbortController and Generic Timer #297
There's a relationship between the queue design and the
EventBussystem, as well as ourWorkerManager.Most important is for us to develop a
Taskabstraction. It can be aclass Task, that represents a "lazy promise". Promises in JS are strictly evaluated, while these tasks will need to be lazily evaluated. Then our task management system can convert our lazy tasks to strict promises (which represents futures). More background info here: https://en.wikipedia.org/wiki/Futures_and_promisesOur task manager will need to have configurable:
- Concurrency limit - indicates the bound on the pool of currently executing tasks (can be 1 to unbounded/infinite)
- Executor - choosing to execute by Node's event-loop, or by passing it into the
WorkerManagerto be executed in a separate thread or core, the former should be used for IO-bound tasks, the latter should be used for CPU-bound tasks. This could be specified by the task creator, rather than the task manager itself.
Stretch goal is to also incorporate "time/calendar-scheduling" so that tasks can be executed at a point in time like cron.
Interaction between
EventBusand task manager may be considered. The event bus is about communicating changes between domains, but the task manager is the one actually executing the tasks.Tasks can be:
- Re-ordered or reprioritised or given priorities
- Can be observed for success or failure
- Can have errors handled
- Can be cancelled using our design for abort signal, and have their real side-effects cancelled
- Can be monitored for progress
- Can be persistent - backed by leveldb
This rabbit hole for this goes deep. So we should make sure not to feature creep our non-blocking task queuing needs.
Example of prior work: https://github.com/fluture-js/Fluture
Also to clarify, we are not creating a "generic distributed job queue", that's the realm of things like redis queue and https://en.wikipedia.org/wiki/List_of_job_scheduler_software. There's so much of this already. We just need something in-process relative to Polykey.
Along with the configurable concurrency limit and executor, I think we should have an interface for the queue as well. Depending on the situation we may need just a simple queue, a priority queue, a persistent database queue like discovery uses, etc etc...
So far as theQueuecares it only needs to supportpushandshift. So we can make theQueuea generic class and pass it any implementation we want for storing the queue so long as it extends the interface.It shouldn't be too hard to make the change. We just need to decide if this degree of control is desired. I can see a need for it though.
Is this a part of #326 ?
Is this a part of #326 ?
Nope, this can be done later.
30 remaining items
The
Queuewill need to work with threadsjs queue too: https://threads.js.org/usage-pool. I'm not sure yet if this means ourWorkerManagerwill need to be changed to work with theQueue, since I don't really want there to be 2 queues. MaybeQueueis for managing the queue persistence, while embedding theWorkerManagerin-memory queue that is one to one for each task that is persisted.Alternatively we actually don't use
WorkerManagerpool, and instead manage our own "pool" directly. This means either extendingPoolfrom threadsjs if possible. However ideally theQueuecan also work without theWorkerManagerbeing available, but I'm not sure if this is possible without very different behaviour. Without there being worker threads, you don't really have the task pulling behaviour, instead one just assigns tasks to a concurrency limit.Note that threadsjs has 2 concurrency limits:
- The parallel number of workers to launch
- The number of concurrent async tasks to run
The only real reason to use workers is to run CPU-intensive tasks, not IO-intensive tasks. (Hard to know precisely until we do benchmarking). So the number of concurrent async tasks should be 1. Therefore we ignore
2.in our design of theQueue.If the worker manager was not available, the parallel number of workers to launch should be the same concurrency limit of the number of tasks to run concurrently in our
Queue. In the former, we would use theos.cpu()count, the latter, this can be specified with some number, with it defaulting to1...Actually we can always default it to
1, and then override it with theos.cpu()count if we expect to supply it with workers.
One issue with this is that the
WorkerManageris also used for other things where extra CPU intensive tasks is just offloaded. Perhaps insteadWorkerManagerstays the same, and the injection of worker manager, means a concurrent number of tasks are dispatched to the worker manager but awaited for normally. We would need some way of checking the capacity of the workers before pushing a task into it.I haven't completed the full design of
Taskclass. But I suspect it needs to be similar to lazy promise here: https://github.com/sindresorhus/p-lazy/blob/main/index.js, and even threadsjs representation uses athenmethod to allowawaitto work on their objects. Their type is:/** * Task that has been `pool.queued()`-ed. */ export interface QueuedTask<ThreadType extends Thread, Return> { /** @private */ id: number; /** @private */ run: TaskRunFunction<ThreadType, Return>; /** * Queued tasks can be cancelled until the pool starts running them on a worker thread. */ cancel(): void; /** * `QueuedTask` is thenable, so you can `await` it. * Resolves when the task has successfully been executed. Rejects if the task fails. */ then: Promise<Return>["then"]; }
If
Taskis in fact aclass Task extends Promise, it would have properties that would be enumerable, and properties that are not. We may need to specify this explicitly: https://debugmode.net/2020/06/18/how-to-make-a-property-non-enumerable-in-javascript/Alternative is to form an a plain object like threadsjs does instead of using classes.
There are some interesting timer APIs: https://nodejs.org/api/timers.html#timeoutrefresh
Note that since
TaskIdis aIdSortable, it's strictly monotonic due to our storing of the last task ID...But this assumes the last Task ID is always stored, and we are intending on deleting tasks off the schedule once completed. I wasn't thinking keeping historical tasks are useful (except for maybe debugging? Although it seems like it would be dropped in production, and logging/tracing systems should be maintaining the audit log).
This means the last task ID may be undefined. So we would store the last Task ID regardless of whether there are any tasks left in the scheduler.
Furthermore, when the clock is shifted backwards, the time will be incremented by 1 until it is greater than the last time. The 1 is the smallest unit of precision, in which case this would be 1 millisecond.
Afterwards, it will be strictly monotonic ID but have a weakly monotonic timestamp up to 4096 IDs per millisecond. After 4096 it would roll over.
The expectation is that it's not possible to generate more than 4096 IDs in a millisecond, so by that time, the time must have increased by at least 1 millisecond.
Anyway this means we need to store
Scheduler/lastTaskIdseparate from theScheduler/taskslevel.Benchmark in js-workers shows that the overhead to call the workers takes about 1.16 to 1.5ms.
Worker Threads Intersection.xlsx
A CPU intensive task should be greater than that time to be worth sending to the worker.
However most scheduling work seems it might not actually be CPU intensive. Like NodeGraph and Discovery is mostly IO. I suppose discovery may have have CPU work to pattern match the data to find the right data on the pages it loads, but this should be dominated by the time spent on IO.
Furthermore sending it to a worker can introduce locking problems. The async locks do not work across the worker threads, they only work within the same event loop context. They are not thread-safe nor process-safe.
This should mean that we should not directly integrate
WorkerManagerinto theQueue, instead individual domains may have their handlers directly pass work to theWorkerManager. TheQueuedoes not decide this since it does not know the nature of the task. The domain that registers the handlers can decide the nature of the task. So they can execute within the main thread, or send it off to a web worker and await for it.This means naturally the
Queuecan have either 1 as a concurrency limit or0to indicate unbound concurrent limit. With an unbound concurrency limit, it just immediately proceeds to execute everything that is due for execution.Priority only comes into play with a concurrency limit so that things get put into priority order. Otherwise all tasks will be asynchronous and immediately executed.
The worker's concurrency/parallel limit is not a concern of the
Queuethen.- linked a pull request that will close this issueFeature TaskManager Scheduler and Queue and Context Decorators for Timed and Cancellable #438
on Aug 8, 2022 - removed a link to a pull requestIntegrate js-db and concurrency testing #419
on Aug 8, 2022 We decided not to bother with preventing resource starvation, however an idea is like this.
- Take advantage of DB's natural key ordering.
- Create a bimap index of Priority/Timestamp -> Task Id AND Timestamp/Priority -> Task Id
- Now we can iterate task ids based on 2 compound indexes: highest priority + earliest timestamp AND earliest timestamp + highest priority
- Use dynamic programming/kinetic priority function that iterates through both sublevels (indexes) simultaneously to fill up a fixed concurrency pool (if unlimited, this policy is unnecessary, just iterate through as fast as posssible)
Simultaneous iteration that uses the timestamp to weight the priority, where the timestamp delta starts from 0 and goes towards infinity. Once could say that this multiples the priority based on a "rate". A delta of 0 multiplies by 1. A delta of infinity multiplies by infinity. Therefore the rate produces a multiplier between 1 to infinity.
Here is an example of the 3 policies:
pM = tD + 1- linearpM = tD^e- exponential (increasing at an increasing rate)pM = ln(tD + 1) + 1- logarithmic (increasing at an decreasing rate)
Once we have the tasks system, all other domains should not have any kind of background processing implemented, they should delegate ALL of that functionality into the tasks system.
The task management is ready. However integration into discovery and nodes domains is being done in #445.
Priority management is static, we won't bother with dynamic priority in #329 (comment) before we see it be a problem.
Issue description here is still relevant to #445, since it contains notes on how best to refactor the discovery system.
Specification
Unattended discovery was added in #320, however, there is no concept of priority within the queue. There are three ways that a vertex (a node or identity) can be added to the discovery queue, and they should follow this order of priority:
queueDiscoveryByNode()andqueueDiscoveryByIdentity()(these are called in the commandsidentities discover(explicit discovery) andidentities trust(explicitly setting a permission, so we want the Gestalt to be updated via discovery).Vertexes with a higher priority should be discovered first, either by being placed at the front of the queue or by modifying the traversal method of the queue. The priority queue could also be further optimised by grouping vertices from the same gestalt together when this is known (for example when adding child vertices).
Additional context
Tasks