Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,11 @@ class SubmissionFileMetadata(TypedDict):
INCIDENTS_COUNTER = Counter(
"kcidb_incidents", "Number of incidents ingested", ["ingester", "origin"]
)
WORKER_FAILURES_COUNTER = Counter(
"kcidb_ingester_worker_failures",
"Number of ingester worker processes that exited abnormally",
["ingester", "reason"],
)


def standardize_tree_names(
Expand Down Expand Up @@ -569,10 +574,25 @@ def ingest_submissions_parallel( # noqa: C901 - orchestrator with IO + multipro
process_queue.qsize(),
)
last_progress = time.time()
if not any(w.is_alive() for w in writers):
Comment thread
tales-aparecida marked this conversation as resolved.
if not process_queue.empty():
logger.error("All workers exited while queue still has items")
break
time.sleep(1)

for writer in writers:
writer.join()
Comment thread
tales-aparecida marked this conversation as resolved.
if writer.exitcode:
reason = "signal" if writer.exitcode < 0 else "exception"
logger.error(
"Worker %s exited with code %s (%s)",
writer.pid,
writer.exitcode,
reason,
)
WORKER_FAILURES_COUNTER.labels(
ingester=INGESTER_GRAFANA_LABEL, reason=reason
).inc()
except KeyboardInterrupt:
out("\nKeyboardInterrupt: terminating workers...")
for writer in writers:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,10 @@

import pytest

from kernelCI_app.constants.ingester import AUTOMATIC_LAB_FIELD
from kernelCI_app.constants.ingester import (
AUTOMATIC_LAB_FIELD,
INGESTER_GRAFANA_LABEL,
)
from kernelCI_app.management.commands.helpers.kcidbng_ingester import (
SubmissionFileMetadata,
_extract_origins_info,
Expand Down Expand Up @@ -835,7 +838,7 @@ def test_ingest_submissions_parallel_success(
mock_value.side_effect = [mock_ok, mock_fail, mock_processed]

# Mock the process
mock_process_instance = MagicMock()
mock_process_instance = MagicMock(exitcode=0)
mock_process.return_value = mock_process_instance

ingest_submissions_parallel(
Expand Down Expand Up @@ -870,3 +873,147 @@ def test_ingest_submissions_parallel_success(
mb / total_elapsed,
),
)

@patch("kernelCI_app.management.commands.helpers.kcidbng_ingester.out", MagicMock())
@patch(
"kernelCI_app.management.commands.helpers.kcidbng_ingester.WORKER_FAILURES_COUNTER"
)
@patch("kernelCI_app.management.commands.helpers.kcidbng_ingester.logger")
@patch("multiprocessing.Process")
@patch("multiprocessing.Queue")
@patch("multiprocessing.Value")
@patch("time.sleep", MagicMock())
@patch("time.time", MagicMock(side_effect=TIME_MOCK))
@patch("os.path.getsize")
def test_ingest_submissions_parallel_worker_failure(
self,
mock_getsize,
mock_value,
mock_queue_cls,
mock_process,
mock_logger,
mock_failures_counter,
):
"""Test that all workers failing breaks the loop and logs errors."""
file1_path = SUBMISSION_FILEPATH_MOCK + SUBMISSION_FILENAME_MOCK

mock_getsize.return_value = self.FILE1_SIZE

mock_queue = MagicMock()
mock_queue_cls.return_value = mock_queue

# Queue not empty on first check, triggers the is_alive guard
mock_queue.empty.return_value = False
mock_queue.qsize.return_value = 1

mock_ok = MagicMock(value=0)
mock_fail = MagicMock(value=1)
mock_processed = MagicMock(value=0)
mock_value.side_effect = [mock_ok, mock_fail, mock_processed]

# Two workers: one crashed (exitcode=1), one killed (exitcode=-9)
crashed_worker = MagicMock(exitcode=1, pid=1001)
crashed_worker.is_alive.return_value = False

killed_worker = MagicMock(exitcode=-9, pid=1002)
killed_worker.is_alive.return_value = False

mock_process.side_effect = [crashed_worker, killed_worker]

ingest_submissions_parallel(
json_files=[file1_path],
tree_names={},
dirs=SUBMISSION_DIRS_MOCK,
max_workers=2,
)

# Loop should break because no workers are alive
mock_logger.error.assert_any_call(
"All workers exited while queue still has items"
)

# Both workers should be joined and their exit codes logged
crashed_worker.join.assert_called_once()
killed_worker.join.assert_called_once()

mock_logger.error.assert_any_call(
"Worker %s exited with code %s (%s)", 1001, 1, "exception"
)
mock_logger.error.assert_any_call(
"Worker %s exited with code %s (%s)", 1002, -9, "signal"
)

# Prometheus counter incremented with correct labels
mock_failures_counter.labels.assert_any_call(
ingester=INGESTER_GRAFANA_LABEL, reason="exception"
)
mock_failures_counter.labels.assert_any_call(
ingester=INGESTER_GRAFANA_LABEL, reason="signal"
)
assert mock_failures_counter.labels.return_value.inc.call_count == 2

@patch("kernelCI_app.management.commands.helpers.kcidbng_ingester.out", MagicMock())
@patch(
"kernelCI_app.management.commands.helpers.kcidbng_ingester.WORKER_FAILURES_COUNTER"
)
@patch("kernelCI_app.management.commands.helpers.kcidbng_ingester.logger")
@patch("multiprocessing.Process")
@patch("multiprocessing.Queue")
@patch("multiprocessing.Value")
@patch("time.sleep", MagicMock())
@patch("time.time", MagicMock(side_effect=TIME_MOCK))
@patch("os.path.getsize")
def test_ingest_submissions_parallel_partial_worker_failure(
self,
mock_getsize,
mock_value,
mock_queue_cls,
mock_process,
mock_logger,
mock_failures_counter,
):
"""Test that partial worker failure still completes when a worker survives."""
file1_path = SUBMISSION_FILEPATH_MOCK + SUBMISSION_FILENAME_MOCK

mock_getsize.return_value = self.FILE1_SIZE

mock_queue = MagicMock()
mock_queue_cls.return_value = mock_queue

# Queue non-empty on first iteration, then empty (surviving worker drained it)
mock_queue.empty.side_effect = [False, True]
mock_queue.qsize.return_value = 1

mock_ok = MagicMock(value=0)
mock_fail = MagicMock(value=0)
mock_processed = MagicMock(value=0)
mock_value.side_effect = [mock_ok, mock_fail, mock_processed]

# One worker crashed, one is still alive and draining the queue
crashed_worker = MagicMock(exitcode=1, pid=1001)
crashed_worker.is_alive.return_value = False

alive_worker = MagicMock(exitcode=0, pid=1002)
alive_worker.is_alive.return_value = True

mock_process.side_effect = [crashed_worker, alive_worker]

ingest_submissions_parallel(
json_files=[file1_path],
tree_names={},
dirs=SUBMISSION_DIRS_MOCK,
max_workers=2,
)

# Both workers should be joined
crashed_worker.join.assert_called_once()
alive_worker.join.assert_called_once()

# Only the crashed worker should trigger an error log and counter
mock_logger.error.assert_called_once_with(
"Worker %s exited with code %s (%s)", 1001, 1, "exception"
)
mock_failures_counter.labels.assert_called_once_with(
ingester=INGESTER_GRAFANA_LABEL, reason="exception"
)
mock_failures_counter.labels.return_value.inc.assert_called_once()
110 changes: 110 additions & 0 deletions docs/ingester.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
# Ingester

The ingester is a Django management command that reads KCIDB submission
JSON files from a spool directory and writes them into the PostgreSQL
database. It runs as a long-lived process on `db.kernelci.org`.

Entry point: `backend/kernelCI_app/management/commands/helpers/kcidbng_ingester.py`

## Parallel ingestion

`ingest_submissions_parallel()` is the main orchestration function.
It uses `multiprocessing.Process` workers and a shared
`multiprocessing.Queue` to parallelize file parsing and database writes.

### Data flow

```
json_files (list)
|
v
[batch into groups of INGEST_FILES_BATCH_SIZE]
|
v
multiprocessing.Queue (maxsize = INGEST_QUEUE_MAXSIZE)
| | |
v v v
worker 0 worker 1 ... worker N-1
(process_batch) (process_batch) (process_batch)
| | |
v v v
DB writes via DB writes via DB writes via
flush_buffers() flush_buffers() flush_buffers()
```

### Worker-queue protocol

1. **Spool phase** - The main process batches `json_files` into groups
of `INGEST_FILES_BATCH_SIZE` and puts each batch into the queue.
This happens before any workers start.

2. **Start workers** - `max_workers` child processes are spawned. Each
runs `process_batch()`, which loops on `process_queue.get()` until
it receives `None`.

3. **Enqueue poison pills** - One `None` per worker is put into the
queue alongside worker starts. Since batches are already in the
queue (FIFO), the Nones always sit behind all batches and workers
consume them only after all real work is done.

4. **Wait for completion** - The main loop waits for the queue to
drain, reporting progress every `progress_every_sec` seconds. If
all workers exit while items remain in the queue, the loop logs an
error and breaks instead of hanging indefinitely.

5. **Join** - After all workers exit, the main process joins them and
prints a final progress report.

### Worker internals (process_batch)

Each worker:

- Closes inherited DB connections (`connections.close_all()`) and opens
fresh ones, since connections cannot be shared across processes.
- Accumulates parsed instances (issues, checkouts, builds, tests,
incidents) into an in-memory buffer.
- Flushes the buffer to the database via `flush_buffers()` whenever
any entity type reaches `INGEST_BATCH_SIZE`.
- Sorts instances by ID before flushing to prevent deadlocks when
multiple workers update the same rows concurrently.
- On exit (receiving `None`), flushes any remaining buffered instances.

### Error handling

- **File parse errors**: The file is moved to the `failed` directory
and the `stat_fail` counter is incremented. The worker continues
processing subsequent files.
- **Worker crash**: `is_alive()` returns False, so the main loop
breaks and joins the remaining workers. After join, non-zero exit
codes are logged and the `kcidb_ingester_worker_failures` counter
is incremented with `reason="exception"` (positive exit code) or
`reason="signal"` (negative, e.g. SIGKILL/-9 from OOM killer).
- **Worker hang**: Depends on `requests` timeouts (set on all
production HTTP calls) to eventually unblock the worker.
- **KeyboardInterrupt**: The main process terminates all live workers
and joins them.

### Log excerpts

When `CONVERT_LOG_EXCERPT` is enabled and `STORAGE_TOKEN` is set,
large log excerpts (exceeding `LOGEXCERPT_THRESHOLD` bytes) are
compressed with gzip and uploaded to external storage. The log excerpt
field is then replaced with a URL reference. An in-memory cache
(`CACHE_LOGS`) deduplicates uploads by SHA-256 hash.

See: `backend/kernelCI_app/management/commands/helpers/log_excerpt_utils.py`

### Relevant constants

Defined in `backend/kernelCI_app/constants/ingester.py`:

- `INGEST_FILES_BATCH_SIZE` - Number of files per queue batch
- `INGEST_BATCH_SIZE` - Number of DB instances before flushing
- `INGEST_QUEUE_MAXSIZE` - Bounded queue size (backpressure)
- `LOGEXCERPT_THRESHOLD` - Byte threshold for uploading log excerpts

Defined in `backend/kernelCI_app/constants/general.py`:

- `REQUESTS_TIMEOUT_UPLOAD_IN_SECONDS` (30s) - Timeout for log excerpt uploads to storage
- `REQUESTS_TIMEOUT_WEBHOOK_IN_SECONDS` (10s) - Timeout for Discord webhook posts
- `REQUESTS_TIMEOUT_FETCH_IN_SECONDS` (30s) - Timeout for fetching external log pages
Loading