Skip to content

Repository files navigation

replicator

Retrieval, fingerprinting, and temporary storage layer for the Cannabis Observer cluster.

Replicator owns content fetching, temporary storage, and fingerprinting for the cluster — the network-bound, byte-handling work re-homed out of Watcher. It is driven by commands on the Redis change bus and reports outcomes as facts:

content.fetch (command)  →  fetch  →  fingerprint  →  temp-store  →  blob_available (fact)
                         ↘  closed without bytes  ──────────────→  fetch_failed  (fact)

content.replicate (cmd)  →  guards  →  create-if-absent  ────────→  replication_complete (fact)
                         ↘  refused / conflict  ────────────────→  replication_failed  (fact)

content.fetch-policy (config)  →  per-host request spacing applied to that fetch

Each command's two outcomes share one stream — content.blobs for a fetch, content.artifacts for a replicate — so an issuer's one consumer group sees either outcome of its command.

content.fetch-policy is the third stream kind and the only one Replicator reads without a consumer group: it carries how often each host may be asked, last-write-wins per host, replayed from the beginning at every boot. The numbers are the issuer's — Replicator enforces spacing, it does not decide it. See docs/contracts/replicator-boundaries.md.

The founding design lives in docs/plans/2026-06-25-replicator-mvp-design.md.

Issuing content.fetch commands? Read docs/contracts/content-fetch-issuer-contract.md first — it is the normative issuer contract and its permanent home, with the refusal list and trust posture in content-fetch-issuer-reference.md and the failure taxonomy in content-fetch-outcome-reference.md. Publish through co-core's to_wire, never hand-rolled fields. The wire carries one domain key — info_source_id, echoed onto both facts and read by nothing here — but correlation is still entirely the issuer's job, on command_id. Most ways of getting either wrong fail silently.

Issuing content.replicate? That loop runs, and it writes for gcs (#29) — read docs/contracts/content-replicate-issuer-contract.md first; its reference carries the reasoning behind the clauses. The contract settled the trust model and the issuer obligations ahead of the code, because a write to a permanent store is bounded by nothing a read is (#34), and the code has since been built to it. Two things bite issuers hardest: the destination is refused, never repaired — you render it, and a redelivery must render the same string, because under T4 that string is the idempotency key — and a command with no fact is a command still being retried, not one that failed.

Issuing content.persist? Read docs/contracts/content-persist-issuer-contract.md. The digest is the address: record it, never a URL, and issue on receipt of the revision, before the temp blob's seven days run out. The loop is off on every host until the broker grants the stream.

Shape

Worker-first. The primary process is a bus consumer (a co-core-aio consumer group on the content.fetch stream), not an HTTP API. A thin FastAPI app exposes /health for local checks; it is not part of the MVP loop and is dev-only until a status surface is wanted.

The Redis broker is cluster infrastructure on its own node (co-broker, operated from CannObserv/broker) — Replicator is a client over the tailnet and does not ship or manage a broker. The >=7.0 server floor is Replicator-critical: AsyncBusConsumer.claim_stale_page reads XAUTOCLAIM's three-element reply, added in Redis server 7.0.

Setup

# Mirror the private cannobserv package index into ./.wheelhouse (co-core, co-core-aio, co-core-sync).
# Requires GOOGLE_APPLICATION_CREDENTIALS (see Environment below).
uv run --no-project --with 'google-cloud-storage>=2,<4' python scripts/sync_wheelhouse.py

uv sync
uv run pre-commit install

Environment

Two env files, with a deliberate boundary between them:

  1. /etc/replicator/.env — production configuration, managed manually on the VM. This is the only file the systemd service reads.
  2. .env (repo root, git-ignored) — dev/agent secrets, e.g. GitHub PATs. Loaded by developers and agents in a shell; never by the service, which has no use for them.

For local work, load both:

set -a; . /etc/replicator/.env 2>/dev/null; . .env 2>/dev/null; set +a

Variables the service uses all live in /etc/replicator/.env. docs/ENVIRONMENT.md is the authoritative reference — every variable, its default, and the reasoning behind each one. Below is only what this deployment actually overrides; anything absent from /etc/replicator/.env runs on the default baked into src/core/config.py.

If you pre-create the blob directory, make it and every parent traversable (0755). The worker sets modes only on directories it creates itself — an existing one keeps whatever mode its operator gave it — so a 0700 level anywhere in the chain leaves every blob_uri unopenable by the service that consumes the fact. Startup logs a warning naming each blocking level.

Blobs are temporary and are reaped. A second task in the worker walks the tree every REPLICATOR_BLOB_SWEEP_INTERVAL_SECONDS, removing blobs untouched for REPLICATOR_BLOB_TTL_SECONDS, .tmp debris older than REPLICATOR_BLOB_TEMP_GRACE_SECONDS, and the shard directories those emptied. The TTL runs from last reference, so re-fetching unchanged bytes restarts it. Disk pressure never shortens it: over REPLICATOR_BLOB_MAX_TOTAL_BYTES the worker stops fetching and leaves commands on the bus rather than deleting bytes a consumer was promised.

Variable Set on this VM Purpose
GOOGLE_APPLICATION_CREDENTIALS /etc/replicator/co-gcs-replicator-writer.json The worker's ADC — the writer SA (co-gcs-replicator-writer@co-gcs, #114), never the wheelhouse reader
REPLICATOR_WHEELHOUSE_CREDENTIALS /etc/replicator/co-pypi-reader.json Read-only key for the wheelhouse mirror, so the boot step never holds the writer
REPLICATOR_REDIS_URL redis://replicator:<password>@broker:6379/0 Change-bus client URL — the replicator ACL user on co-broker, over the tailnet
REPLICATOR_BLOB_BACKEND gcs Temp blobs live in an object store — not the local default
REPLICATOR_BLOB_BUCKET co-gcs-blobs The temp-blob bucket that backend writes into
REPLICATOR_PERSIST_ENABLED true, since 2026-09-26 Whether this worker consumes content.persist (#114). Default off; turned on here after the broker granted the stream (broker#64). Needs REPLICATOR_PERMANENT_BUCKET
REPLICATOR_PERMANENT_BUCKET co-gcs-replicator, set at the #114 step 5 deploy The permanent content-addressed store: a second replicate source, and where content.persist writes (#114); unset ⇒ none
REPLICATOR_BLOB_DIR /var/lib/replicator/blobs Temp-storage root; unused under gcs, kept against a flip back to local
REPLICATOR_REPLICATION_ALIASES_FILE /etc/replicator/replication-aliases.json The alias table: gcs-publication → gs://co-gcs-publication with the publication writer's credentials_file, empty prefix — since the #114 cutover (2026-09-25). primary bound the same bucket until it was retired on 2026-09-27, and gs://co-gcs-replication before that (#86). root:exedev 640, like the key files beside it. Unset ⇒ nothing provisioned, every content.replicate command refused alias_unknown
REPLICATOR_CONSUMER_NAME (unset) Per-group override; the name is derived from the group — replicator-fetch-1 — and never shared
REPLICATOR_REPLICATE_CONSUMER_NAME (unset) The same, for replicator.replicate — derives replicator-replicate-1
REPLICATOR_PERSIST_CONSUMER_NAME (unset) The same, for replicator.persist — derives replicator-persist-1
REPLICATOR_LOG_LEVEL INFO Root log level

BUILD_ID is stamped by the unit's ExecStartPre rather than set in the env file. Every other REPLICATOR_* setting — the blob TTL and ceilings, the consumer group and start id, the read window and pacing fallback, the reclaim and backoff numbers — is on its default; see docs/ENVIRONMENT.md for what each one is and why.

Seeding a fetch

Watcher has issued content.fetch since watcher#241, so the deployed loop is proven by its traffic, not by a seeded command: each stored a blob and published blob_available in sudo journalctl -u replicator is one of Watcher's commands closing.

scripts/seed_fetch.py publishes to scratch streams. The target is never defaulted — --redis-url and --topic are both required — and --dry-run prints the frames without contacting a broker at all:

# The scratch redis-server the integration tests use (docs/TESTING.md) — reaches no worker.
uv run python -m scripts.seed_fetch \
  --redis-url redis://localhost:6379/15 --topic replicator.itest.seed \
  https://example.test/a

The live stream is not an example here (#90). db 0 + content.fetch still needs --production and a real --info-source-id, but a frame there is fetched for real on a command Watcher never issued — an operator act under Watcher's identity, never this host's replicator credential, which reaches that stream only through an ACL gap CannObserv/broker#14 closes. Every flag, and what --watch can see: docs/COMMANDS.md.

Test & lint

uv run pytest                          # default suite; integration and gcs deselected
uv run pytest --no-cov -m integration  # a scratch redis-server, never the broker, plus the
                                       # OOM rows against a broker they spawn (#79)
uv run pytest --no-cov -m gcs          # T4 rows against the GCS test bucket (#38)
uv run ruff check .

Full command reference: docs/COMMANDS.md.

Dev server

The FastAPI /health app runs on port 8001 (port 8000 belongs to systemd if/when the API is promoted to a deployed surface):

uv run uvicorn src.api.main:app --host 0.0.0.0 --port 8001 --reload --log-config src/core/log_config.json

--log-config routes uvicorn's own uvicorn / uvicorn.access / uvicorn.error loggers — which ship with propagate=False and plain-text handlers of their own — through the same JSON formatter configure_logging() installs on the root logger. Without it the dev server emits mixed-format output: plain-text access lines interleaved with JSON app records.

Deploy

The systemd unit lives at deploy/replicator.service and runs the worker, not the API. To install on a fresh host:

# Copy into systemd's path — all four, they are copies and not symlinks
sudo cp deploy/replicator.service /etc/systemd/system/replicator.service
sudo cp 'deploy/replicator-failure-notify@.service' /etc/systemd/system/
# tailscaled's -900, and the slice grant without which MemoryLow= is inert (#113)
sudo install -D -m 644 deploy/tailscaled.service.d/memory.conf /etc/systemd/system/tailscaled.service.d/memory.conf
sudo install -D -m 644 deploy/system.slice.d/replicator-memory.conf /etc/systemd/system/system.slice.d/replicator-memory.conf
sudo systemctl daemon-reload
sudo systemctl restart tailscaled   # OOMScoreAdjust= applies at exec only
sudo systemctl enable --now replicator

# The host's memory tunables, which sysctl reads — not systemd
sudo cp deploy/99-co-replicator-memory.conf /etc/sysctl.d/
sudo sysctl --system

# apt's needrestart hook lists restarts, never performs them (#122)
sudo install -D -m 644 deploy/needrestart.conf.d/replicator.conf /etc/needrestart/conf.d/replicator.conf

# Proof it took: the live checks read /proc and /sys/fs/cgroup, where
# `systemctl show` would report values the running processes may not have
uv run pytest --no-cov tests/test_deploy.py

# Tail logs
sudo journalctl -u replicator -f

The OnFailure= handler is easy to skip and silent when missing: nothing runs it until something fails, which is the one moment it was supposed to help. What it reported is read with journalctl -t replicator-failure, never -u (docs/FAILURE-NOTIFICATION.md).

Production secrets live in /etc/replicator/.env (managed manually on the VM, not in the repo). The unit's ExecStartPre writes the current git SHA to /run/replicator/build-id and exposes it as BUILD_ID, asserts the Redis >=7.0 floor via scripts/check_redis_floor.sh, and refreshes the wheelhouse via scripts/sync_wheelhouse.py — whose journald output is plain text, not JSON, by design (see docs/STYLE.md, "Not everything in the journal is JSON").

Because ExecStart runs --frozen --no-sync, run uv sync --frozen as part of the deploy, before systemctl restart.

ExecStartPre also refuses the start outright if this checkout is not on main, or if main carries commits that were never pushed (#37, #48) — the deployed code has to be code someone else can see. So get level with origin/main first: git pull --ff-only when the merge happened on GitHub, git push when it happened here. Full verdict table and the escape hatch: docs/DEPLOYMENT.md.

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages