Skip to content

feat: push line ingest and tail streaming jobs (#13) - #41

Merged
cursor[bot] merged 4 commits into
mainfrom
cursor/push-tail-ingestion-1cc6
Aug 17, 2026
Merged

cursor[bot] merged 4 commits into
mainfrom
cursor/push-tail-ingestion-1cc6

Conversation

@leo-aa88

Copy link
Copy Markdown
Member

Closes #13.

Adds direct line push and long-lived tail ingestion so callers do not have to re-upload batch files for every analysis.

Push

POST /v1/ingestions/lines (unversioned alias too) accepts NDJSON of raw lines or JSON objects (message / text / raw). Lines go through the existing parse → fingerprint → persist path. Creates IngestionJob(mode=push). Over INGEST_PUSH_MAX_LINES (default 5000) → 400.

Tail

POST /v1/ingestions with "mode": "tail" (cloudwatch / datadog / loki). One long-lived job; the worker ticks SourceAdapter.read from the saved cursor so new lines stay on the same ingestion_job_id.

Lifecycle:

  • POST /v1/ingestions/{id}:pause
  • POST /v1/ingestions/{id}:resume
  • POST /v1/ingestions/{id}:stop (terminal)

N consecutive tail errors (TAIL_ERROR_THRESHOLD, default 5) auto-pause. /health includes tail_jobs: {running, paused}.

Backpressure

Pending worker queue at INGEST_QUEUE_MAX (default 100) → 429 with Retry-After and { "error_code": "INGEST_QUEUE_FULL" }. This is a stand-in for full G9.

Schema

Alembic 0005_ingest_modes: mode, cursor, last_polled_at, consecutive_errors on ingestion_jobs. Existing rows backfill mode=batch.

Out of scope

G6 content dedup / Idempotency-Key — overlapping tail windows can double-count until that lands. Full G9 API/LLM rate limits are not in this PR.

Open in Web Open in Cursor 

cursoragent and others added 2 commits August 17, 2026 07:21
POST /v1/ingestions/lines accepts NDJSON; tail mode plus pause/resume/stop
are worker-managed. Queue-full requests return 429. Refs #13.

Co-authored-by: Leonardo <leo-aa88@users.noreply.github.com>
Regenerate clients/openapi.json after adding /v1/ingestions/lines and
tail lifecycle paths. Refs #13.

Co-authored-by: Leonardo <leo-aa88@users.noreply.github.com>
@cursor

cursor Bot commented Aug 17, 2026

Copy link
Copy Markdown
Contributor

Bugbot is not enabled for your account, so this pull request was not reviewed.

Enable Bugbot in the Cursor dashboard to get automatic reviews on future PRs.

@cursor

cursor Bot commented Aug 17, 2026

Copy link
Copy Markdown
Contributor

Code review — PR #41 (issue #13)

Reviewed origin/main...origin/cursor/push-tail-ingestion-1cc6 (commits 1a7f381, 39a727a). CI green noted. Applied Alembic 0001–0004 are untouched; 0005_ingest_modes adds mode, cursor, last_polled_at, consecutive_errors. Push goes through _process_line (parse → fingerprint → persist). /lines is a static route (GET 405, not UUID 400). Dual-mount /v1 + deprecated aliases, ingest/admin on POST lines and :pause/:resume/:stop (version-stripped), 429 + Retry-After + INGEST_QUEUE_FULL, pause/resume/stop 409, auto-pause, and /health.tail_jobs all check out. G6/G9 full implementations correctly left out of scope.

must-fix

None

should-fix

  • src/core/ingestion/tail.py:171 + src/worker/runner.py:186-189 (tick_one_tail_job failure path / run_worker sleep). A failed tick does not update last_polled_at, so _due_tail_jobs still selects the job (last_polled_at IS NULL or still older than TAIL_POLL_INTERVAL). tick_tail_jobs then returns len(jobs) > 0, and the worker skips its sleep. Concrete failure: one adapter blip is retried in a tight loop until TAIL_ERROR_THRESHOLD (default 5) auto-pauses — no 30s spacing, and the remote API is hammered. If a tick flushed a 500-line batch then raised, those retries re-read the same window immediately (stronger duplicate burst than the documented G6 adjacent-window overlap). Fix: either always sleep the poll interval after a loop, or treat failed ticks as not-due until last_polled_at/next_retry_at advances.

  • src/core/ingestion/tail.py:196-198 (tick_one_tail_job success path). On success, last_polled_at is always set to now even when ingest_from_source persisted a non-null pagination cursor (_store_job_cursors). Datadog read() keeps ref.cursor when max_rows is hit and sends that cursor with the request from/to. The next tick builds TimeWindow(start=last_polled_at, end=now) and reapplies the old page cursor. Concrete failure: a 30s window with > DATADOG_MAX_ROWS (default 10000) events drops the remainder of that window and can page the new range with a stale cursor. Loki is fine (cursor is a timestamp and overrides window.start); CloudWatch usually drains the window so cursor is None. Fix: do not advance last_polled_at while any stream cursor is still set, or keep the previous window until the cursor is exhausted.

nice-to-have

  • tick_tail_jobs / _due_tail_jobs: SELECT … FOR UPDATE SKIP LOCKED loads all due tails then holds those row locks across adapter I/O. POST :pause on a locked job waits until every due tick in that transaction finishes. Claim/tick/commit one job at a time.
  • GET /v1/ingestions still filters status == "completed", so running/paused/stopped tails never appear in the UI picker (detail-by-id still works).
  • TestTickAutoPause.test_success_resets_counter mocks ingest_from_source and does not assert resume_cursors=…, existing_job=job, finalize=False. test_health_includes_tail_job_counts uses one mocked scalar_one() for queue depth, running, and paused, so it cannot tell those queries apart.
  • OpenAPI for POST /v1/ingestions/lines documents only 200; 400 (over cap / empty) and 429 are missing. Lifecycle routes omit 404/409.
  • TAIL_ADAPTERS is a hard-coded name set in core; a registry flag would scale better when a new pull adapter is added.

Verdict

must-fix count: 0
should-fix count: 2

Failed ticks stamp last_polled_at so jobs wait TAIL_POLL_INTERVAL.
Open pagination cursors keep the original window until exhausted. Refs #13.

Co-authored-by: Leonardo <leo-aa88@users.noreply.github.com>
@cursor

cursor Bot commented Aug 17, 2026

Copy link
Copy Markdown
Contributor

Code review — PR #41 (issue #13), round 2

Reviewed origin/main...origin/cursor/push-tail-ingestion-1cc6 after af2e18e (1a7f381, 39a727a, af2e18e). CI green noted. Re-checked the two round-1 should-fix items against the live code, not just the new tests.

Previous should-fix status

  • item 1: fixed — tick_one_tail_job now stamps last_polled_at = now on the failure path (tail.py:261). _due_tail_jobs SQL still uses last_polled_at IS NULL OR last_polled_at <= cutoff, and the Python is_tail_job_due re-filter matches that predicate. After a failed tick the worker may skip sleep for one loop (ticked == 1), then the job is not due and idle-sleeps until TAIL_POLL_INTERVAL. test_failed_tick_is_not_due_until_poll_interval drives the real tick_one_tail_job + _due_tail_jobs path; MagicMock ignores SQL WHERE, but production SQL would also drop the job. Not a false pass.
  • item 2: fixed (success path) — _apply_tail_window_progress keeps last_polled_at at window.start and stores metadata_json["tail_window_end"] while any stream cursor is non-null; tick_window_for_job reuses that exact from/to. test_open_cursor_holds_window_until_exhausted patches tick_window_for_job only on the first tick, then calls the real function and asserts the hold, then exhausts the cursor and asserts last_polled_at advances and tail_window_end is cleared. Matches Datadog max_rows keeping ref.cursor bound to the request from/to.

must-fix

None

should-fix

  • src/core/ingestion/tail.py:168-171 + src/core/ingestion/tail.py:261 (_apply_tail_window_progress / failure path). The two round-1 fixes overload last_polled_at: hold uses it as window start, backoff uses it as last attempt. A continuation tick that raises (Datadog AdapterUnavailableError with 0 new lines on page 2 — ingest_from_source re-raises when every stream fails before yielding) sets last_polled_at = now and leaves the open cursor + tail_window_end in place. tick_window_for_job then fails the hold (last_polled_at < stored_end is false once now is past the original end, which it usually is) and builds TimeWindow(start=failure_time, end=now) with the stale page token. Concrete failure: remainder of the max_rows-split window is dropped; Datadog may 400 or page the wrong range. The new test never fails a tick while a cursor is open, so this path is untested. Fix: keep window bounds in metadata (tail_window_start / tail_window_end) and always use last_polled_at as the poll clock, including on failure, so backoff and hold no longer share a field.

nice-to-have

None new (round-1 nice-to-haves still apply; not re-raised).

Verdict

must-fix count: 0
should-fix count: 1

Store tail_window_start/end in metadata; last_polled_at is poll backoff
only so a failed continuation tick cannot drop an open page. Refs #13.

Co-authored-by: Leonardo <leo-aa88@users.noreply.github.com>
@cursor

cursor Bot commented Aug 17, 2026

Copy link
Copy Markdown
Contributor

Code review — PR #41 (issue #13), round 3

Reviewed origin/main...origin/cursor/push-tail-ingestion-1cc6 after 7f34204 (1a7f381, 39a727a, af2e18e, 7f34204). CI green noted. Re-checked the round-2 field-split should-fix against the live code and the new test, not just the commit message.

Previous should-fix status

  • round-2 field split: fixed — last_polled_at is always the poll clock (now on success in _apply_tail_window_progress and on failure in tick_one_tail_job); while a cursor is open, tick_window_for_job returns _held_paging_window() from metadata_json['tail_window_start'|'tail_window_end'] and never uses last_polled_at as start. test_failed_continuation_keeps_held_window actually fails a tick while a cursor is open: first tick persists page-2 + window bounds, second tick raises, asserts last_polled_at == fail_now and backoff, then tick_window_for_job(job, later) is the original 12:00–12:01 rather than fail_now → later.

must-fix

None

should-fix

None

nice-to-have

None (round-1 nice-to-haves still apply; not re-raised)

Verdict

must-fix count: 0
should-fix count: 0

@cursor
cursor Bot merged commit cc8bcc4 into main Aug 17, 2026
2 checks passed
@leo-aa88
leo-aa88 deleted the cursor/push-tail-ingestion-1cc6 branch August 17, 2026 07:50
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

feat(G4): push (POST /v1/ingestions/lines) + tail streaming ingestion

2 participants