# Events

## Produced Streams

| Stream | Event type | Consumer |
| --- | --- | --- |
| `repo.snapshot.ready` | `RepoSnapshotReadyEvent` | Code Parser Service |
| `files.changed` | `FileChangedEvent` | Code Parser Service, Docs Ingestion Service |
| `commits.changed` | `CommitsChangedEvent` | Code Parser Service |

Stream name constants live in `internal/contracts/events/streams.go`.

## Consumed Streams

None — repo sync is a producer-only service.

## Contract Source

Event definitions are replicated locally under `internal/contracts/` (no shared `code-intelligence-contracts` module dependency for now):

| Package | Purpose |
| --- | --- |
| `internal/contracts/metadata/` | v1 envelope (`event_id`, `event_version`, `created_at`) |
| `internal/contracts/events/` | Event structs, constructors, `ChangeType` and `FileKind` enums |
| `internal/contracts/validate/` | Pre-publish validation |
| `internal/contracts/errors/` | Validation error types |

## Event envelope

Every published message includes:

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

Use the `New*Event()` constructors in `internal/contracts/events/` to populate the envelope automatically.

## Downstream routing

Both Code Parser and Docs Ingestion consume the same `files.changed` stream and filter by `file_kind`:

| `file_kind` | Consumer | Purpose |
| --- | --- | --- |
| `code` | Code Parser | Parse source files into symbols/chunks |
| `docs` | Docs Ingestion | Parse Markdown/docs into structured artifacts |
| `config`, `test`, `asset`, `unknown` | Future consumers | Reserved for later routing rules |

Docs Ingestion publishes `docs.parsed.ready` (future stream). Code Parser consumes parsed docs output and combines it with code parse results when building `graph.artifact.ready`.

## Required fields

### repo.snapshot.ready

- `repo_id`, `snapshot_id`, `commit_sha`, `ref`

Optional: `artifact_uri`, `file_count`

### files.changed

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

Optional: `artifact_uri`, `old_path` (for renames)

`file_kind` values: `code`, `docs`, `config`, `test`, `asset`, `unknown`

`change_type` values: `added`, `modified`, `deleted`, `renamed`

#### Example: code file

```json
{
  "event_id": "evt_01",
  "event_version": "v1",
  "created_at": "2026-05-25T12:00:00Z",
  "repo_id": "AD/example-repo",
  "snapshot_id": "snap_01",
  "commit_sha": "abc123",
  "file_path": "internal/config/config.go",
  "language": "go",
  "file_kind": "code",
  "change_type": "added"
}
```

#### Example: docs file

```json
{
  "event_id": "evt_02",
  "event_version": "v1",
  "created_at": "2026-05-25T12:00:01Z",
  "repo_id": "AD/example-repo",
  "snapshot_id": "snap_01",
  "commit_sha": "abc123",
  "file_path": "docs/architecture.md",
  "language": "markdown",
  "file_kind": "docs",
  "change_type": "added"
}
```

### commits.changed

- `repo_id`, `snapshot_id`, `commit_sha`

Optional: `parent_sha`, `ref`, `author`, `message`, `changed_files`

## Future stream: docs.parsed.ready

Planned producer: Docs Ingestion Service. Planned consumer: Code Parser (graph composition stage).

Required fields (planned):

- `repo_id`, `snapshot_id`, `commit_sha`, `file_path`, `doc_id`, `artifact_uri`, `doc_type`

Optional: `title`, `section_count`

Example payload:

```json
{
  "event_id": "evt_03",
  "event_version": "v1",
  "created_at": "2026-05-25T12:00:02Z",
  "repo_id": "AD/example-repo",
  "snapshot_id": "snap_01",
  "commit_sha": "abc123",
  "file_path": "docs/architecture.md",
  "doc_id": "doc_01",
  "artifact_uri": "file:///data/artifacts/doc_01.json",
  "doc_type": "markdown",
  "title": "Architecture",
  "section_count": 8
}
```

Go structs for this stream are not implemented yet.

## Publishing mechanism

Repo Sync publishes indexing events with Redis **`XADD`**, not BullMQ or other job-queue libraries.

```text
validate(event) → json.Marshal → XADD <stream> * payload <json>
```

| Step | Detail |
| --- | --- |
| Validate | `internal/contracts/validate/` rejects incomplete events before publish |
| Stream name | Base name from `internal/contracts/events/streams.go`, optional prefix via `REDIS_STREAM_PREFIX` |
| Payload field | Single Redis hash field `payload` containing the full event JSON |
| Idempotency | Stable `event_id` in the envelope; consumers dedupe on retry |

Example (conceptual):

```text
XADD files.changed * payload {"event_id":"evt_01","event_version":"v1",...}
```

Downstream services (Code Parser, Docs Ingestion, Commit Intelligence) consume with **`XREADGROUP`** and a dedicated consumer group per service. See [stream-topology.md](../../code-intelligence-contracts/docs/stream-topology.md) in the shared contracts repo.

### Redis Streams vs BullMQ

| | Redis Streams | BullMQ |
| --- | --- | --- |
| Purpose | Event bus (“something happened”) | Job queue (“do this task”) |
| Publish API | `XADD` | `queue.add()` (Node.js) |
| Multi-consumer | Yes — same stream, consumer groups | Typically one worker per job |
| This pipeline | Yes | Not used |

## Validation before publish

Call the matching validator in `internal/contracts/validate/` before writing to Redis Streams:

- `validate.RepoSnapshotReady(event)`
- `validate.FileChanged(event)`
- `validate.CommitsChanged(event)`

## Publish workflow

On first clone or sync (orchestrated in `SyncService.SyncRepository`):

1. Git sync completes → persist repository + snapshot metadata
2. Create `indexing_runs` record with expected counts; bulk-insert `snapshot_files` inventory
3. Publish one `repo.snapshot.ready`
4. Publish one `files.changed` event per discovered/changed file (all files are `added` on first clone)
5. Publish one `commits.changed` event per commit in chronological order (full history on first sync; new commits only on re-index)
6. Mark indexing run `repo_sync` stage completed

Downstream consumers:

- Code Parser processes `files.changed` where `file_kind=code`
- Docs Ingestion processes `files.changed` where `file_kind=docs`
- Code Parser (graph deltas) consumes `commits.changed` for ordered delta computation

Implementation:

- [`internal/redis/publisher.go`](../internal/redis/publisher.go) — `EventPublisher` using `XADD`
- [`internal/service/sync.go`](../internal/service/sync.go) — orchestration workflow
- [`internal/service/indexing_run.go`](../internal/service/indexing_run.go) — indexing run + snapshot file inventory
- [`internal/service/publish.go`](../internal/service/publish.go) — maps `git.SyncResult` to contract events with deterministic `event_id` values:
  - `evt_{snapshot_id}_ready`
  - `evt_{snapshot_id}_file_{hash8}`
  - `evt_{snapshot_id}_commit_{hash8}`
