Skip to content

Generic Non-Blocking Task Management ("Queue") for discovery and nodes domains #329

Description

@emmacasolin

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:

  1. Manually, via the discovery methods queueDiscoveryByNode() and queueDiscoveryByIdentity() (these are called in the commands identities discover (explicit discovery) and identities trust (explicitly setting a permission, so we want the Gestalt to be updated via discovery).
  2. As a step in the discovery process whereby child vertices are added into the discovery queue in order to discover the entire connected gestalt.
  3. Automatically by a process of rediscovery when we want to update existing Gestalts in the Gestalt Graph (to be addressed in Discovery - revisiting Gestalt Vertices and error handling #328).

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

  1. Modify the existing Discovery Queue to be a Priority Queue
  2. Ensure that when a user interactively wants to discover a gestalt vertex that it becomes the highest priority and gets executed first
  3. Look into the potential for further optimising the priority queue, for example by having multiple points of comparison with varying levels of importance that can influence the priority of a particular vertex in the queue

Activity

  1. CMCDragonkai commented on Feb 15, 2022

    @CMCDragonkai
    Member

    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 on IdSortable second.

    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.

  2. CMCDragonkai commented on Feb 16, 2022

    @CMCDragonkai
    Member

    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.

  3. CMCDragonkai commented on Feb 23, 2022

    @CMCDragonkai
    Member

    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.

  4. CMCDragonkai commented on Apr 20, 2022

    @CMCDragonkai
    Member

    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.

  5. emmacasolin commented on Apr 21, 2022

    @emmacasolin
    ContributorAuthor

    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.

  6. 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
  7. CMCDragonkai commented on Apr 26, 2022

    @CMCDragonkai
    Member

    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:

    There's a relationship between the queue design and the EventBus system, as well as our WorkerManager.

    Most important is for us to develop a Task abstraction. It can be a class 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_promises

    Our 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 WorkerManager to 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 EventBus and 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.

  8. CMCDragonkai commented on Apr 26, 2022

    @CMCDragonkai
    Member
  9. CMCDragonkai commented on Apr 26, 2022

    @CMCDragonkai
    Member

    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.

  10. tegefaulkes commented on Apr 27, 2022

    @tegefaulkes
    Contributor

    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 the Queue cares it only needs to support push and shift. So we can make the Queue a 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.

  11. tegefaulkes commented on Apr 27, 2022

    @tegefaulkes
    Contributor

    Is this a part of #326 ?

  12. CMCDragonkai commented on Apr 27, 2022

    @CMCDragonkai
    Member

    Is this a part of #326 ?

    Nope, this can be done later.

  13. 30 remaining items

  14. CMCDragonkai commented on Aug 8, 2022

    @CMCDragonkai
    Member

    The Queue will need to work with threadsjs queue too: https://threads.js.org/usage-pool. I'm not sure yet if this means our WorkerManager will need to be changed to work with the Queue, since I don't really want there to be 2 queues. Maybe Queue is for managing the queue persistence, while embedding the WorkerManager in-memory queue that is one to one for each task that is persisted.

    Alternatively we actually don't use WorkerManager pool, and instead manage our own "pool" directly. This means either extending Pool from threadsjs if possible. However ideally the Queue can also work without the WorkerManager being 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:

    1. The parallel number of workers to launch
    2. 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 the Queue.

    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 the os.cpu() count, the latter, this can be specified with some number, with it defaulting to 1...

    Actually we can always default it to 1, and then override it with the os.cpu() count if we expect to supply it with workers.


    One issue with this is that the WorkerManager is also used for other things where extra CPU intensive tasks is just offloaded. Perhaps instead WorkerManager stays 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.

  15. CMCDragonkai commented on Aug 8, 2022

    @CMCDragonkai
    Member

    I haven't completed the full design of Task class. 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 a then method to allow await to 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"];
    }
  16. CMCDragonkai commented on Aug 8, 2022

    @CMCDragonkai
    Member

    If Task is in fact a class 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.

  17. CMCDragonkai commented on Aug 8, 2022

    @CMCDragonkai
    Member

    There are some interesting timer APIs: https://nodejs.org/api/timers.html#timeoutrefresh

  18. CMCDragonkai commented on Aug 8, 2022

    @CMCDragonkai
    Member

    Note that since TaskId is a IdSortable, 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/lastTaskId separate from the Scheduler/tasks level.

  19. CMCDragonkai commented on Aug 8, 2022

    @CMCDragonkai
    Member

    Benchmark in js-workers shows that the overhead to call the workers takes about 1.16 to 1.5ms.

    https://github.com/MatrixAI/js-workers/blob/27fb8c3051f0880b81a14ea7daee47b15a94dd89/benches/results/worker_manager_metrics.txt#L2

    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 WorkerManager into the Queue, instead individual domains may have their handlers directly pass work to the WorkerManager. The Queue does 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.

  20. CMCDragonkai commented on Aug 8, 2022

    @CMCDragonkai
    Member

    This means naturally the Queue can have either 1 as a concurrency limit or 0 to 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 Queue then.

  21. CMCDragonkai commented on Aug 23, 2022

    @CMCDragonkai
    Member

    We decided not to bother with preventing resource starvation, however an idea is like this.

    1. Take advantage of DB's natural key ordering.
    2. Create a bimap index of Priority/Timestamp -> Task Id AND Timestamp/Priority -> Task Id
    3. Now we can iterate task ids based on 2 compound indexes: highest priority + earliest timestamp AND earliest timestamp + highest priority
    4. 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:

    desmos-graph (2)

    • pM = tD + 1 - linear
    • pM = tD^e - exponential (increasing at an increasing rate)
    • pM = ln(tD + 1) + 1 - logarithmic (increasing at an decreasing rate)

    https://www.desmos.com/calculator/wnezpfgxqc

  22. CMCDragonkai commented on Sep 11, 2022

    @CMCDragonkai
    Member

    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.

  23. CMCDragonkai commented on Sep 13, 2022

    @CMCDragonkai
    Member

    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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions