Incremental sync
Do one full copy, then apply only what changed.
Incremental sync means copying a workspace once and then applying only what changed, for as long as the integration runs. Four rules make it correct. Get any of them wrong and the sync looks fine while it quietly loses records.
The shape
- Take an event watermark.
- Copy everything by paging the list endpoints.
- Read the change feed from the watermark, from then on.
- Use
updated_afteronly to recover, never as your main loop.
Rule 1: take the watermark before the backfill
The feed is read oldest first, so the watermark is the sequence of the newest event you can see. Read to the end of the feed to find it.
import os
import requests
BASE = "https://tahoe.workonward.com/api/partner/v1"
session = requests.Session()
session.headers["Authorization"] = f"Bearer {os.environ['TAHOE_API_KEY']}"
def current_sequence():
"""The sequence of the newest event this key can see, or None if none."""
after, last = None, None
while True:
params = {"limit": 100}
if after is not None:
params["after"] = after
response = session.get(f"{BASE}/events", params=params, timeout=30)
response.raise_for_status()
body = response.json()
if body["data"]:
last = body["data"][-1]["sequence"]
if not body["has_more"]:
return last
after = body["next_after"]
def cold_start(store):
# 1. The watermark FIRST, held in memory for now.
watermark = current_sequence()
# 2. The backfill. Slow, and that is fine: it happens once.
backfill_jobs(store)
backfill_applicants(store)
backfill_applications(store)
# 3. Save the watermark only now. A crash during the backfill then
# restarts the whole cold start, instead of resuming from a point
# the copy never reached.
store.save_sequence(watermark)Rule 2: the sequence number is the state
Do not track progress with a timestamp. Events are written by changes that happen at the same time and can become visible slightly out of clock order, so “everything since 10:04” can miss an event stamped 10:03 that became visible at 10:05. The sequence number orders events exactly and has no such gap.
Rule 3: save the watermark after the work
def drain(store):
after = store.load_sequence()
while True:
params = {"limit": 100}
if after is not None:
params["after"] = after
response = session.get(f"{BASE}/events", params=params, timeout=30)
response.raise_for_status()
body = response.json()
for event in body["data"]:
# At-least-once delivery: the same event can arrive more than
# once, so this check is not an optimization.
if not store.already_processed(event["id"]):
dispatch(store, event)
store.mark_processed(event["id"])
# AFTER the event is handled. Saving first turns a crash halfway
# through a page into events that are skipped for good. Use the
# event's own sequence: next_after is null on the last page.
after = event["sequence"]
store.save_sequence(after)
if not body["has_more"]:
returnSaving first and working second loses events when you crash. Working first and saving second replays them, and a replay is free when your handlers are safe to run twice. They must be anyway, because delivery is at least once even without crashes.
Rule 4: re-read the resource, do not trust the payload
Event payloads are thin: identifiers and a few plain values. Read the resource by its handle when you handle the event.
| If you use | Then |
|---|---|
| The event payload as your data | You keep what your scopes allowed when the event was written. A paywall change or an erasure later leaves an outdated copy behind. |
| A fresh read by handle | You get what your scopes allow now, and a deleted resource answers 404, which is the right answer. |
One part of the payload is worth using directly. On some events, data.object.changed lists the fields that changed. When none of them is a field you copy, you can skip the read.
# This system records outcomes only, not pipeline stages.
MIRRORED_APPLICATION_FIELDS = {"status"}
def on_application_event(store, event):
changed = event["data"]["object"].get("changed")
# When the event lists what changed and you copy none of it, skip the
# read. When there is no list, read again to be safe.
if changed and not set(changed) & MIRRORED_APPLICATION_FIELDS:
return
handle = event["data"]["object"]["id"]
response = session.get(f"{BASE}/applications/{handle}", timeout=30)
if response.status_code == 404:
store.delete_application(handle) # gone since the event was written
return
response.raise_for_status()
store.upsert_application(response.json())When to use updated_after
The jobs, applications, applicants, sourced profiles and shared pool lists accept ?updated_after=. It is for recovery, not for your steady state:
- You suspect drift and want to check a window of time.
- You added a field to your own schema and need to fill it in.
- You want a backstop. A daily check over the last day or two catches anything the feed did not bring you.
If your reader falls behind
Events are kept for 30 days. A reader that was down for a few days just drains the feed from its watermark. A reader that was down for longer than 30 days has lost events it can never get back. Catch up on erasure notices first, then run the cold start again, and remove any record the new backfill did not return: its deletion event may be one of the events you missed.
Webhooks, and why you still keep the feed
Webhooks push the same events in the same envelope, so one handler serves both. Treat them as a way to hear sooner, not as a replacement.
Deliveries fail, endpoints are disabled after repeated failures, and a delivery that runs out of attempts is given up. The ordered feed is the source of truth and the only way to recover. Keep the drain loop, run it on a schedule as a backstop, and keep saving the watermark from it.
Pacing
- Use
limit=100. One request instead of four. - One worker per key. Parallel workers share the same limits, and the
429answers start to look random. - Honor
Retry-After. For the personal-data budget it can be hours, and it is the value to trust. - Watch
personal_data_reads_used_todayin GET /me during a backfill, and stop on purpose before the limit. Running into it halfway leaves you guessing what was written. - Log the request ID (the
Tahoe-Request-Idheader) with your own run record, so support can find the exact request.