# Events

Code Parser consumes Repo Sync indexing streams and publishes parser output streams. All messages use Redis Streams with a single `payload` field containing JSON.

Stream name constants live in `internal/contracts/events/streams.go`. Optional prefix via `REDIS_STREAM_PREFIX` (e.g. `staging` → `staging.files.changed`).

## Consumed Streams

| Stream | Event type | Handler |
| --- | --- | --- |
| `files.changed` | `FileChangedEvent` | Validate → filter `file_kind=code` → parse file → persist artifact → publish outputs |
| `repo.snapshot.ready` | `RepoSnapshotReadyEvent` | Validate → register snapshot → fan-in check → maybe finalize merged graph |
| `commits.changed` | `CommitsChangedEvent` | Validate → resolve commit diff → load base/target artifacts → compute graph delta → save delta → publish `graph.delta.ready` |

Consumer group: `code-parser-service` (override with `REDIS_STREAM_GROUP`). Strict retry: ACK on durable success or permanent skip only; transient failures stay pending; `XAUTOCLAIM` reclaims stale messages; after max retries (`REDIS_STREAM_MAX_RETRIES`, default 5) messages move to `{stream}.dlq`.

### files.changed routing

Code Parser processes only `file_kind=code` events. Other kinds are ACK'd without parsing.

| `file_kind` | Action |
| --- | --- |
| `code` | Parse via language analyzer |
| `docs`, `config`, `test`, `asset`, `unknown` | Skip (ACK only) |

Skip rules also apply for `change_type=deleted` and unsupported languages.

## Produced Streams

| Stream | Event type | When published |
| --- | --- | --- |
| `chunks.ready` | `ChunksReadyEvent` | After successful parse when chunks are non-empty |
| `graph.artifact.ready` | `GraphArtifactReadyEvent` | Once per snapshot after code + docs fan-in; `artifact_uri` is `mongo://snapshot_graphs/art_{snapshot_id}_merged` |
| `parser.diagnostics.ready` | `ParserDiagnosticsReadyEvent` | Parse/skip/publish failures with diagnostics |
| `graph.delta.ready` | `GraphDeltaReadyEvent` | After commit-level graph delta saved to store |

Large deltas are persisted as a manifest in `graph_deltas` plus rows in `graph_delta_parts` when the payload exceeds MongoDB limits. The published `delta_uri` remains `mongo://graph_deltas/{delta_id}`.

Publisher: `internal/redis/publisher.go` — validates, JSON-serializes, `XADD` with `payload` field. Logs `publish_duration_ms` on success.

## Contract Source

Local packages under `internal/contracts/` (wire-compatible with Repo Sync):

| Package | Purpose |
| --- | --- |
| `internal/contracts/metadata/` | v1 envelope (`event_id`, `event_version`, `created_at`) |
| `internal/contracts/events/` | Event structs, constructors, enums |
| `internal/contracts/graph/` | Graph IR (`Artifact`, `Node`, `Edge`) |
| `internal/contracts/validate/` | Pre-publish validation |
| `internal/contracts/errors/` | Validation error types |

Use `New*Event()` constructors to populate the envelope automatically.

## Event envelope

Every message includes:

```json
{
  "event_id": "evt_...",
  "event_version": "v1",
  "created_at": "2026-05-22T11:30:00Z"
}
```

## Pipeline (happy path)

```text
Repo Sync
  └─ files.changed (file_kind=code)
       └─ Code Parser consumer
            ├─ read source from GIT_WORKSPACE_PATH
            ├─ Go analyzer → graph IR + chunks
            ├─ upsert artifact → MongoDB `code_graph_artifacts`
            ├─ upsert chunks → MongoDB `code_chunks` (text hydrated from source lines)
            ├─ increment `indexing_runs.processed_code_files`
            ├─ chunks.ready (when chunks non-empty)
            └─ (no per-file graph.artifact.ready)

Fan-in (code files + docs complete):
            ├─ merge per-file artifacts + doc_chunks/sections
            ├─ upsert `snapshot_graphs`
            └─ graph.artifact.ready (one per snapshot)

Repo Sync (or future publisher)
  └─ commits.changed
       └─ Code Parser consumer
            ├─ resolve changed code files (event or git diff)
            ├─ load base/target graph artifacts
            ├─ compute delta → save delta JSON (inline or chunked in `graph_deltas` + `graph_delta_parts`)
            └─ graph.delta.ready
```

## Required fields

### files.changed (consumed)

- `repo_id`, `snapshot_id`, `commit_sha`, `file_path`, `language`, `file_kind`, `change_type`

### graph.artifact.ready (produced)

- `repo_id`, `snapshot_id`, `commit_sha`, `artifact_id`, `artifact_uri`, `parser_version`, `graph_schema_version`
- Optional counts: `node_count`, `edge_count`

### chunks.ready (produced)

- `repo_id`, `snapshot_id`, `commit_sha`, `file_path`, `language`, `chunks[]`
- Each chunk: `chunk_id`, `symbol_name`, `symbol_type`, `start_line`, `end_line`

### parser.diagnostics.ready (produced)

- `repo_id`, `snapshot_id`, `commit_sha`, `job_id`, `diagnostics[]`
- Each diagnostic: `code`, `message`, `severity`, optional `file_path`, `line`

## Example: files.changed → outputs

**Input** (`files.changed`):

```json
{
  "event_id": "evt_files_01",
  "event_version": "v1",
  "repo_id": "ad/example",
  "snapshot_id": "snap_01",
  "commit_sha": "abc123",
  "file_path": "internal/config/config.go",
  "language": "go",
  "file_kind": "code",
  "change_type": "modified"
}
```

**Output** (`graph.artifact.ready`):

```json
{
  "event_id": "evt_snap_01_graph_...",
  "event_version": "v1",
  "repo_id": "ad/example",
  "snapshot_id": "snap_01",
  "commit_sha": "abc123",
  "artifact_id": "art_snap_01_deadbeef",
  "artifact_uri": "file:///data/artifacts/art_snap_01_deadbeef.json",
  "parser_version": "v0.1.0",
  "graph_schema_version": "v1",
  "node_count": 8,
  "edge_count": 12
}
```

## Debug REST APIs

Inspect in-memory job state and on-disk artifacts (see [`api.md`](api.md)):

- `GET /parse-jobs/{job_id}`
- `GET /artifacts/{artifact_id}`
- `GET /diagnostics/{job_id}`
