DEV Community

Robin
Robin

Posted on

Make Fragment Epoch an Invariant Before a Restarted Scorer Publishes a Grade

I still replay a grading-lane fixture I built for this review, because a green dashboard can hide a lying grade. A worker wrote fragment one, the eval server restarted, and a new worker wrote another fragment one. The delayed fragment two from the first epoch then arrived, and the publisher appended it anyway. Would you trust that mixed score inside a canary gate, or would you halt the lane?

The common implementation treats each partial score as an append-only log keyed only by attempt id. That key alone is not enough, because a restart creates a new generation that still shares the attempt. The invariant I want is narrower than idempotent writes: a grade may contain fragments from only one epoch. If two epochs mix, should the publisher reject the late fragment, or should it compensate the grade?

Assumptions I will not quietly expand

I am reviewing an architecture here, and I am not claiming a measured production outage or latency. Assume one attempt id, many fragments, at-least-once delivery, and a scorer that can restart mid-run. Assume the free model call is remote, slow enough to overlap a restart, and not a local deterministic function. Assume you will re-check current product terms, because I will not invent hardware, duration, or a permanence promise.

Disclosure: This article was prepared as part of MonkeyCode's product outreach. I would use this open-source project's free model access and free server as the fixture's remote scoring plane. The outreach brief states a free allowance of ten million tokens, which I did not measure myself. Have you confirmed that figure in current docs before sizing fan-out, or are you planning from a screenshot?

Constraints, then the data flow

I am constraining this review to one attempt stream, one publisher, and one remote scoring worker at a time. Cross-attempt fairness, semantic model quality, and multi-region consensus all stay outside the scope of this review. A restart may reorder fragments, duplicate a sequence number, or drop the worker without a clean cancel. The publisher is the only component allowed to make a finished grade visible to the canary gate.

The attempt store mints an epoch, the queue carries that epoch on every fragment, and the worker must echo it. The free server may execute the model call, but it must not choose a new epoch after a crash. A late fragment with a stale epoch is a different failure domain from a duplicate inside the live epoch. Do you see why a write keyed only by attempt id collapses those two domains into one silent append?

sequenceDiagram
    participant Store as Attempt store
    participant Queue as Fragment queue
    participant Worker as Scoring worker
    participant Fence as Epoch fence
    participant Gate as Canary gate
    Store->>Queue: enqueue(attempt, epoch=1, seq=0)
    Worker->>Fence: fragment(epoch=1, seq=0)
    Note over Worker: process restarts
    Store->>Queue: bump epoch=2 and re-enqueue
    Worker->>Fence: fragment(epoch=2, seq=0)
    Queue->>Fence: late fragment(epoch=1, seq=1)
    Fence-->>Gate: reject stale epoch, do not publish mix

Numbered protocol I want in front of publish

I want the protocol to be boring, because boring fences are the ones reviewers can actually test. A fence that lives only in a diagram will not survive the first restarted worker you inject. So I put the check in the publisher, and I also put the epoch on the queue message itself. Why keep the check in both places, if a single publisher check already covers the happy path?

  1. Mint a new epoch when the attempt is admitted, and persist that epoch before the first model call is enqueued.
  2. Stamp the attempt id, the epoch, and the sequence number on every fragment the worker emits, including retries.
  3. Reject any arriving fragment whose epoch does not equal the open epoch currently stored for that attempt.
  4. Treat a repeated epoch and sequence pair as a duplicate acknowledgement, not as a second sample to append.
  5. Publish only when the expected sequence set is complete and the set of epochs inside the grade has size one.
  6. On restart, bump the stored epoch and drop every unpublished fragment that still belongs to the previous epoch.
  7. Surface rejects on a side log so the halt gate can see loss, instead of hiding that loss inside a retry counter.

Failure domains the fence does not magically erase

I split the lane into five failure domains, because a single retry policy will smear them together. The model plane can time out after tokens have already been spent, and that waste is not the same thing as a wrong epoch. The worker process can restart or get preempted, which is the domain this fence is actually aimed at. The queue can reorder or duplicate, and the gate can still read a grade if you publish before the fence runs.

There is a nasty overlap I do not want you to miss when you draw the next box. A zombie worker from epoch one can still be inside a model call after epoch two is opened. The fence can reject its late fragment, but it cannot unspend the call that already left the process. Is that a correctness bug in the grade, or a cost leak I should not hide inside the publisher check?

A fixture you can run before you trust the plane

This next block is an unexecuted local fixture, not a benchmark I collected from a shared server. I want you to run it on your laptop before you point the same state machine at any remote worker. The asserts encode the property, and a green printout only means the fixture held, not that a vendor plane is durable. If a line fails, fix the fence before you debate which model should sit behind the worker.

#!/usr/bin/env python3
'''Local epoch-fence fixture. Not a product benchmark.'''

from dataclasses import dataclass, field


@dataclass(frozen=True)
class Fragment:
    attempt_id: str
    epoch: int
    seq: int
    body: str


@dataclass
class GradeBook:
    open_epoch: dict[str, int] = field(default_factory=dict)
    rows: dict[str, list[Fragment]] = field(default_factory=dict)
    rejected: list[Fragment] = field(default_factory=list)

    def begin(self, attempt_id: str, epoch: int) -> None:
        self.open_epoch[attempt_id] = epoch
        self.rows.setdefault(attempt_id, [])

    def restart(self, attempt_id: str) -> int:
        nxt = self.open_epoch[attempt_id] + 1
        self.rows[attempt_id] = []
        self.open_epoch[attempt_id] = nxt
        return nxt

    def append(self, fragment: Fragment) -> str:
        current = self.open_epoch.get(fragment.attempt_id)
        if current is None or fragment.epoch != current:
            self.rejected.append(fragment)
            return 'reject'
        seen = {(row.epoch, row.seq) for row in self.rows[fragment.attempt_id]}
        if (fragment.epoch, fragment.seq) in seen:
            return 'duplicate'
        self.rows[fragment.attempt_id].append(fragment)
        return 'accept'

    def publish(self, attempt_id: str, expected: int) -> str:
        got = self.rows.get(attempt_id, [])
        epochs = {row.epoch for row in got}
        seqs = {row.seq for row in got}
        if epochs != {self.open_epoch[attempt_id]}:
            return 'halt'
        if seqs != set(range(expected)):
            return 'incomplete'
        return 'publish'


def counterexample() -> None:
    book = GradeBook()
    book.begin('a1', 1)
    assert book.append(Fragment('a1', 1, 0, 'old-sample')) == 'accept'
    assert book.restart('a1') == 2
    assert book.append(Fragment('a1', 2, 0, 'new-sample')) == 'accept'
    assert book.append(Fragment('a1', 1, 1, 'late-old')) == 'reject'
    assert book.publish('a1', expected=1) == 'incomplete'
    assert book.append(Fragment('a1', 2, 1, 'new-tail')) == 'accept'
    assert book.publish('a1', expected=2) == 'publish'
    epochs = {row.epoch for row in book.rows['a1']}
    assert epochs == {2}, epochs
    print(
        'property=single_epoch status=publish rejected=1 '
        'denominator=published_grades numerator=mixed_grades value=0'
    )


if __name__ == '__main__':
    counterexample()
Enter fullscreen mode Exit fullscreen mode
python3 epoch_fence.py
python3 -m py_compile epoch_fence.py
Enter fullscreen mode Exit fullscreen mode

I expect that command to print a single property line, and I expect the late fragment to land in rejected rather than in rows. If your wrapper shells out to a remote worker, keep this process as the publisher and let the worker only propose fragments. Would you let the worker publish directly just because the server was free to start, or would you keep the fence local? A free start is only an availability convenience, and it is not an assignment of epoch ownership.

Injected failures, denominator, and the acceptance rule

I inject four orders before I talk about scale, because scale will only hide the mix. First, deliver sequence zero, restart, deliver a new sequence zero, then deliver the old sequence one late. Second, deliver the same epoch and sequence twice, and demand a duplicate ack instead of a second row. Third, publish before the expected sequence set is complete, and demand incomplete rather than a short grade.

Fourth, leave a zombie fragment in flight across the bump, and demand a reject even if its body looks newer. The denominator is the count of grades that reached a publish status inside this injected suite. The numerator is how many of those published grades contain more than one fragment epoch inside them. Incomplete attempts stay out of the denominator, or a crash will look like a perfect score.

The acceptance rule is a numerator of zero, with every stale-epoch fragment present in the reject log. I would fail the review if a mixed grade is published even once inside that injected suite. I would also fail it if a stale fragment is missing from the reject log after the late delivery. What good is a green gate if the denominator quietly dropped the only attempt that mixed?

Tradeoffs I would write on the review

Choice Latency effect Correctness Cost effect What still breaks
Reject a stale epoch One side-log write Keeps a single-epoch grade Drops already spent calls Zombie spend after the bump
Append by attempt id only Fewest checks Mixes two samples Looks cheap until the gate lies Restart plus reorder
Compensate after publish Extra rewrite Works only if the gate has not read Double write on the grade A reader that already consumed the mix
Wait for quiescence before publish Adds a heartbeat wait Fewer late mixes Holds a worker slot longer A server that never sends the final beat

I would take the reject row for this lane, and I would delay compensation until a later review. Replay is the wrong default here, because replaying a stale fragment is how the mix gets built. Does that feel strict when the late body happens to look better than the new one? It should feel strict, because looks better is not an epoch check you can defend in review.

What I would change next

The next change is a fencing token the worker must present on the model call, not only on the fragment. Without that token, epoch two can be open while epoch one is still spending a remote call. I would also split the reject log from the retry counter, so backpressure is a visible depth and not a hidden loop. I would not add a second publisher until this fixture stays at numerator zero under reorder.

Who should leave this pattern alone

Do not use this fence if your side effect is a payment, an email, or any ledger this grade book does not own. Do not use it if you cannot persist the epoch before enqueue, because a minted-but-lost epoch recreates the mix. Do not treat a stated token allowance as a throughput SLO, and do not skip the local fixture. Teams that need semantic judging of model output still need a separate eval, since this protocol only protects fragment membership.

The question I want left open

A remote worker can host the restart injection if you confirm the current free-model and free-server offer first. I would keep the epoch fence in your process either way, and I would halt publish when the numerator moves off zero. I am not asking you to trust a vendor grade, and I am asking you to trust a property you can rerun. Which concrete event order breaks that single-epoch invariant, and should the system reject, replay, or compensate?

Top comments (0)