Every scheduler writes down the next run
When a weekly report finishes, the scheduler has to record next Monday's run. The obvious place for it is the task, as a next-run timestamp on the definition that is advanced every time the task fires. On one machine, with one process owning the timer, this is a good design. The state is a file that can be opened in an editor when something looks wrong, and backing it up means copying the file.
The design strains when tasks belong to many tenants and run on machines that appear and vanish. Several dispatchers may then reach for the same task at once, and the field that says when it runs next has to be protected from all of them.
Storing the next run on the task has three limitations
Firstly, advancing the series mutates the definition on every tick, and two dispatchers racing to mutate it need a lock, or a claim on the task that expires after a set time. Secondly, such a claim is a lease, which trades a bound on waiting for a guess about how long the work takes[1]. A machine that dies in the middle of a run holds the task until the guess expires, and the choice is between a lease long enough to be safe and one short enough not to hurt. Finally, the definition mixes what a person asked for with the state of its runs. An edit and a tick write to the same row, and nothing separates the fields a person chose from the fields the scheduler maintains.
A definition is intent, and each run is a row
We propose that the definition holds what a person asked for and one switch, its name, description, schedule, repeat rule and entrypoint, and whether it is enabled. It carries no next-run timestamp, no flag saying that it is firing and no lease. Every field on it is authored, and a field derived from a run appearing there counts as a bug. Everything about a particular firing lives in a ledger of occurrences, one row each, and the open row is the thing that will fire. What runs next is then a query and not a field, and safety comes from the key of each row instead of a lease.
Two places to keep the next run. On the left, the task holds its next run and a lease, and two dispatchers race for it, in coral. On the right, the definition holds only intent, each occurrence is a row in a ledger, and both dispatchers compute the same key and adopt the same open row, in blue.
The key is derived, not allocated
The key of an occurrence is composed from the facts that identify it. These are how it is delivered, what wakes it, the agent it belongs to, where it runs, the task, a digest of the task's authored fields and the time it is due. Two components that agree on the facts produce the same string. When each tries to create the occurrence, the second finds the row the first created and adopts it, and the two converge on one row without talking to each other. Nobody holds a lock, and nobody trusts a clock to say when a claim has lapsed.
The key of one occurrence, split into the facts it is made from. The component that projects the occurrence and the dispatcher that fires it each compose the key from what they know, and both arrive at the same row, in blue.
The digest covers authored fields only. A column added to the definition for some other reason changes nothing, while a real edit, such as a new schedule or a different entrypoint, retires the open occurrence and creates a fresh one under a new key.
Three ways it went wrong
Our first version let the runtime advance the series. The component that started an occurrence computed the next slot and wrote it onto the definition. It felt natural, since that component was already holding the task. When definitions stopped carrying run state, the runtime instead projected the next occurrence onto the ledger, and its own comment records the danger, that a series which stops projecting silently stops recurring. Over about a month we found several ways for that write not to happen. In every one the run succeeded and the task never woke again. Projection now belongs to the component that owns the repeat rule. Marking a run as running projects its successor, before the run does any work, and no dispatcher has a handoff it can drop.
The second mistake was in the key itself, which has to be derived identically wherever it is derived, and for a while was not. A destination such as team:11 and a timestamp with an explicit offset were normalised differently by the two implementations. The dispatcher failed to recognise the row projected for it and created a second one. Two rows for one occurrence look exactly like two runs at once, and the guard against overlapping runs then declined to start anything, on every tick, without an error. Both implementations are now pinned to one set of cases that each must reproduce byte for byte.
The third was a sweep we added underneath the whole design, a periodic pass that gives any enabled task without an open occurrence a new one. It ran on time every fifteen minutes and reported success every time, and for two days it repaired nothing. One tenant hit a lock timeout, the whole pass shared one database transaction, and every tenant after it failed inside a transaction that was already dead. From outside, the pass was indistinguishable from a healthy fleet, since a healthy fleet also has nothing to repair.
The three failures, each beside what it looked like from outside. A dropped successor looked like a task nobody had scheduled, a duplicate row looked like a run already in progress, and a dead sweep looked like a fleet with nothing to fix. In each case, in coral, failing looked the same as succeeding.
Related work
Gray and Cheriton introduced leases, locks that expire, which bound how long a failed holder can block others at the cost of a guess about duration[1]. Helland argues that any operation which may be retried must be idempotent, with a repetition indistinguishable from a single attempt[2]. Kubernetes names each job that a scheduled job creates after the schedule and the scheduled time. Its documentation still warns that a job may be created twice or not at all, and that jobs should be idempotent as a result[3]. Of the schedulers we know, the next run is usually kept on the job, with a lease to protect it. We keep it in a ledger of occurrences, where a key derived from the occurrence's own facts makes creating it idempotent, and the definition holds only what was asked for.
Open questions
Firstly, the design has far more machinery than a file and a timer. Nobody can inspect the state by opening a file, and for a single tenant on one machine we would not choose it. Secondly, the sweep walks every project to find the few that have tasks, and its cost grows with the number of tenants, not with the number of tasks. Finally, all three failures stayed hidden because failing looked exactly like succeeding. We have no good answer for that yet. One option is to make the healthy case report something positive. Silence would then stop being the normal state and become evidence of a fault.
However, we argue that the ledger remains the right place for the next run. Each of these concerns how the ledger is operated and observed, and making its silence informative is the one we most want future work to answer.
1. Gray and Cheriton. Leases: an efficient fault-tolerant mechanism for distributed file cache consistency. SOSP 1989.
2. Helland. Idempotence is not a medical condition. ACM Queue, 2012.
3. The Kubernetes authors. CronJob. Kubernetes documentation.