DEV Community

Casey Li
Casey Li

Posted on

Commit the Offset Before the Draft

A consumer watermark is a durability fact, not a suggestion. Free model access can draft fixtures and review notes. It must not choose, delay, or advance an offset commit. The process that read the record is the only one allowed to stamp progress.

Picture a warehouse ledger and the rubber stamp beside it. A clerk may pencil a note about a strange crate, and that note can even be useful at the next shift change. The stamp that marks the crate as received stays with the person who counted the crate. Helpful prose does not move inventory.

An offset is that stamp. Once it moves forward, the log treats earlier records as done, and a later restart will not hand them back. Duplicates are often recoverable through an idempotent handler. Silent skips are not, because nothing in the next poll announces the hole.

The commit therefore stays on the client even when a team wants a second opinion on the policy. A remote draft can be slow, empty, or confidently wrong. Waiting on that draft before committing leaves a window where a rebalance retires the assignment epoch the code still believes it holds. Applying a suggested higher offset deletes the gap between the true mark and the fiction from the retry set.

Some teams still want a scratch environment for sample batches and for a written critique of the commit policy. Disclosure: This article was prepared as part of MonkeyCode's product outreach. MonkeyCode offers free model access and a free server option, and either can hold that scratch work: fixture files, a policy diff, a review comment. Neither option is a commit log, a consumer group, or a fence.

This article does not name models, token quotas, hardware, duration, or permanence. Those details were not verified for this draft, and inventing them would turn a boundary guide into an advertisement with fake specifications. A design that needs those numbers before the consumer can move is already outside the approach described here.

A notebook that shows a draft matching a fixture set is encouraging, and it is still only a comment on examples. That result can justify keeping an offline review step on a scratch workspace. It cannot justify moving the stamp, any more than a passing lint run justifies skipping the commit protocol.

Red flags show up as ordinary code, not as a dramatic outage banner. The first is a parser. Production code that reads a model string and casts the last token to an integer has already handed the watermark away, regardless of how careful the prompt sounded.

Shared membership is the next flag. A free server process that joins the same consumer group, or a scratch disk treated as the source of truth for the committed offset, has promoted an optional seat into a replica.

The third flag is narrower than a general debate about retries. A failed commit call whose backoff is chosen by prompting, instead of by the client library inside the process that holds the assignment, has moved a protocol decision off the client. The fourth flag is a test that treats a model score as proof the offset is safe. A score is not a fence epoch.

Better alternatives are dull, which is the point. Keep a checkpoint next to the consumer, and bind every commit to the epoch issued when the assignment was granted. Advance the mark only after the handler returns success and the epoch still matches. Generate fixtures as plain files, review them as data, and merge policy changes through the same path as any other patch.

The broker client, or a local file that stands in for it during tests, performs the commit. A free server seat, when one already exists, is a reasonable scratch workspace for the fixture directory so it does not clutter a laptop. It is a poor place to park a group coordinator, a member id, or the watermark itself.

The modules below are a proposal. They were not executed while this article was written. An operator should run the commands before treating the pattern as local evidence, and should keep the files out of any process that talks to a real broker until that run is green.

# checkpoint_gate.py — proposal, unexecuted in this draft
from __future__ import annotations

import json
from dataclasses import dataclass
from pathlib import Path


@dataclass(frozen=True)
class Assignment:
    group: str
    partition: int
    epoch: int


class Checkpoint:
    def __init__(self, path: Path) -> None:
        self.path = path

    def load(self) -> dict:
        if not self.path.exists():
            return {"epoch": None, "offset": -1}
        return json.loads(self.path.read_text(encoding="utf-8"))

    def commit(self, assignment: Assignment, offset: int) -> None:
        current = self.load()
        stored = current["epoch"]
        if stored not in (None, assignment.epoch):
            raise RuntimeError(
                f"fence mismatch: stored={stored} offered={assignment.epoch}"
            )
        if offset < current["offset"]:
            raise RuntimeError("refusing to move the watermark backward")
        payload = {
            "epoch": assignment.epoch,
            "offset": offset,
            "group": assignment.group,
            "partition": assignment.partition,
        }
        tmp = self.path.with_suffix(".tmp")
        tmp.write_text(json.dumps(payload), encoding="utf-8")
        tmp.replace(self.path)


def apply_advisory(advisory: dict, assignment: Assignment, store: Checkpoint) -> None:
    """Review comments never move the watermark, even if they name an offset."""
    del advisory, assignment, store
Enter fullscreen mode Exit fullscreen mode

The anti-pattern is short enough to keep beside the gate as a warning. It is not a function the consumer loop should import, even behind a feature flag.

# do_not_ship.py — illustrated failure, not an API
def commit_from_model_text(text: str, consumer, partition: int) -> None:
    offset = int(text.strip().split()[-1])
    consumer.commit({partition: offset})
Enter fullscreen mode Exit fullscreen mode

That function fails in three independent ways. It trusts prose, it carries no epoch, and it can skip every record the model did not mention. Deleting it is the fix. Wrapping the same parse in a helper with a nicer name is not a fix.

A pair of tests locks the boundary. The advisory names a much higher offset, and the stored mark must stay put. A stale epoch must be rejected even when the offered offset looks like the next honest step.

# test_checkpoint_gate.py — proposal, unexecuted in this draft
from pathlib import Path

from checkpoint_gate import Assignment, Checkpoint, apply_advisory


def test_advisory_cannot_advance(tmp_path: Path) -> None:
    store = Checkpoint(tmp_path / "wm.json")
    live = Assignment("invoices", 3, epoch=7)
    store.commit(live, 10)
    apply_advisory(
        {"suggested_offset": 99, "note": "draft from a free model"},
        live,
        store,
    )
    loaded = store.load()
    assert loaded["offset"] == 10
    assert loaded["epoch"] == 7


def test_stale_epoch_is_refused(tmp_path: Path) -> None:
    store = Checkpoint(tmp_path / "wm.json")
    store.commit(Assignment("invoices", 3, epoch=7), 10)
    try:
        store.commit(Assignment("invoices", 3, epoch=6), 11)
    except RuntimeError as exc:
        assert "fence mismatch" in str(exc)
    else:
        raise AssertionError("stale epoch was allowed to commit")
Enter fullscreen mode Exit fullscreen mode

Fixture generation stays boring on purpose. It writes JSON and returns. It does not open a socket, and it does not import the checkpoint type. If that script grows a commit call, review stops before any discussion of prompt quality.

# write_fixtures.py — proposal, unexecuted in this draft
import json
from pathlib import Path


def write_fixtures(out: Path, count: int = 24) -> None:
    out.mkdir(parents=True, exist_ok=True)
    rows = [
        {"id": i, "amount_cents": 100 + i, "kind": "invoice.posted"}
        for i in range(count)
    ]
    target = out / "batch.json"
    target.write_text(json.dumps(rows, indent=2) + "\n", encoding="utf-8")


if __name__ == "__main__":
    write_fixtures(Path("scratch/batch"))
Enter fullscreen mode Exit fullscreen mode

The local check is a short command sequence from a clean virtual environment. The fragment assumes a POSIX shell, and a Windows host needs the equivalent activation step before the same pytest invocation. No credential is required, because the gate has no vendor client and the fixture writer has no network call.

python -m venv .venv
. .venv/bin/activate
python -m pip install pytest
python -m pytest test_checkpoint_gate.py -q
python write_fixtures.py
python -c "import pathlib; text = pathlib.Path('checkpoint_gate.py').read_text(); assert 'socket' not in text and 'commit_from_model_text' not in text"
Enter fullscreen mode Exit fullscreen mode

The last command is a crude tripwire. It will not catch a dynamic import hidden inside a constructed string, so a human still reads the diff. Its job is to fail fast when a network call or the anti-pattern name is pasted into the file that is supposed to be a stamp.

The table is a field guide, not a scorecard. The commit column never says to ask a model. The exit column is the point at which this whole split should be abandoned rather than patched with a cleverer prompt.

Situation Where progress is stamped What a free draft may do Exit if this appears
Handler succeeded and the epoch matches Client checkpoint or broker commit Comment on policy text, offline Draft text is parsed into an offset
Epoch changed after a rebalance Nobody, until the new owner commits Nothing on the hot path Scratch host still stores the old member id
Fixture batch for review Nowhere Write JSON in a scratch workspace Fixture script imports the consumer
Commit call failed Client-library backoff in the assigned process An incident note after the fact A prompt is asked whether to retry the commit
A quota or uptime figure is required to proceed Not this design Not a substitute for that figure The consumer blocks on the free option

Exit criteria are operational. Drop model involvement on this path when any one of them is true.

The consumer cannot make progress unless the free option is reachable. An auditor requires every commit input to be replayed without a third party.

Someone wants to hide a draft hop inside the handler because the latency budget looks tight on a diagram. A reviewer cannot explain the fence without opening a prompt. The scratch workspace starts holding group membership, offsets, or credentials.

The replacement in each of those cases is the same shape. Use a broker-native commit, keep a tested local fence in the assigned process, and limit model output to documents that never ship. Removing the model from the path entirely is a valid outcome of the field guide, not a failure of it.

Several teams should not adopt the split at all. A pipeline whose only runtime is a free server seat should not pretend that seat is a consumer group. Payment, clinical, and legal logs, where a skipped offset is an integrity incident, should not collect policy advice on the same path that stamps progress. Operators who need a published quota, a named model, or an uptime commitment before accepting a dependency will not find those facts in this article, and should not infer them.

Anyone who would auto-apply a generated patch onto Checkpoint.commit is already past the first red flag. The useful response is to stop the merge, not to add a second model that checks the first.

Limitations sit next to the proposal so they are harder to skip. The samples remain unexecuted until the commands above are run on a machine the operator controls. They imitate the shape of an offset commit. They are not a Kafka client, and they do not speak a broker protocol.

Filesystem replace is not a replicated log. A crash between the temp write and the replace can leave a stray file, which a real commit protocol handles with more care than this stand-in. The gate does not encrypt the checkpoint, authenticate the caller, or address poison records, idempotent handlers, or dead-letter routing.

Those neighboring problems stay neighboring on purpose. Moving any of them onto free inference is a different design, and this article does not bless that move. No latency, cost, or accuracy figure here should be quoted elsewhere, because none was measured.

The habit that survives contact with a rebalance is small. Stamp progress in the process that holds the assignment. Let a draft comment from outside the fence. When the draft and the stamp trade places, put the watermark back on the client and leave it there.

Top comments (0)