An agent on every host — at-least-once log shipping and offset checkpoints

A shipping agent on each host owns three pieces of state — a byte offset on disk (positions.yaml / registry), an in-RAM batch, and a retry slot — and ships batches at-least-once, which means duplicates on retry are normal and a checkpoint that advances before the ship can silently lose lines.

Previously

Centralising means putting a network between the app and its log file — and that something on each host that has to keep the pipe full has a name. It is the shipping agent: Filebeat, Promtail, Fluent Bit, or Vector. They look different on the outside; under the skin they are the same three compartments.

Scene 02

An agent on every host

  1. Watch
  2. Try it
  3. Predict
  4. Capture
HAPPY PATH · acks flowingHOST/var/log/app.logoffset=02025-05-08T12:00 app: ...2025-05-08T12:01 app: ...2025-05-08T12:02 app: ...2025-05-08T12:03 app: ...← tail cursorshipping agent1 · POSITIONS.YAML/var/log/app.logoffset: 02 · BATCH BUFFER0/43 · RETRY SLOT(empty)BACKENDreceiver:3100received 0 linesPOST /loki/api/v1/pushidleHappy path. The agent tails app.log, batches lines, POSTs, waits for the ack, then advances positions.yaml — …
What to watch for

The app writes lines into /var/log/app.log. The agent reads each new line into its in-memory batch. When the batch hits 4 lines, it POSTs the batch over HTTP. Only AFTER the backend acks does positions.yaml advance — that on-disk offset is what the agent re-reads from on restart.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

Agent.tail_loop
reads new lines from app.log into the in-RAM batch
def tail_loop():
f = open('/var/log/app.log')
f.seek(self.offset) # offset from positions.yaml
while running:
line = f.readline()
if not line:
sleep(poll_interval); continue
batch.append(line)
if len(batch) >= batch_capacity:
send_queue.put(batch)
batch = []
Agent.send_loop
POSTs batches; advances positions only after a 2xx ack
def send_loop():
while running:
batch = retry_slot or send_queue.get()
resp = http.post(backend_url, batch)
if 200 <= resp.status < 300:
self.offset += sum(len(l) for l in batch)
persist(positions_yaml, self.offset)
retry_slot = None
else: # network drop, 5xx, or no ack
retry_slot = batch # resend after backoff
sleep(backoff_with_jitter())
Agent.on_restart
the only durable state is positions.yaml — RAM is gone
def on_restart():
# batch and retry_slot lived in RAM — both gone
self.offset = read(positions_yaml) or 0
# Promtail gotcha: if positions advanced when the line
# was READ rather than when the backend ACKED, every
# in-flight line between offset and EOF is silently lost.
spawn(tail_loop) # re-opens app.log, seeks to offset
spawn(send_loop) # starts with empty queue + slot

Where this sits in Build a distributed logging stack (ELK / Loki)

Scene 02 of 12. A shipping agent tails each log file from a saved offset, batches lines, and POSTs them at-least-once — duplicates on retry are normal, not a bug.

Up next. The agent ships at-least-once over a pipe — but ships them WHERE? The backend at the far end can stall, and when it does the agent has exactly three options for what to do with the lines piling up — only one of which keeps the application alive.

All 12 scenes in Build a distributed logging stack (ELK / Loki) · Every curriculum

Built with Arqly
Every scene in Build a distributed logging stack (ELK / Loki) builds on the one before it.All 12 Build a distributed logging stack (ELK / Loki) scenes