Commit 6791d73
authored
feat(server): refactor server ingestion to sink (#165)
## Summary
Refactored the server observability write path to use the shared
`ControlEventSink` abstraction while preserving the existing
`/api/v1/observability/events` behavior and Postgres-backed OSS storage.
Added a default OSS server sink that adapts the existing `EventStore`
write path into the shared sink contract.
Kept the current `DirectEventIngestor` usage stable by allowing existing
store-based construction to continue working, while routing writes
through sink semantics internally.
## Scope
### User-facing / API changes
- No intended user-facing or HTTP API changes.
- `/api/v1/observability/events` request/response behavior remains
unchanged.
- `IngestResult` and API response accounting semantics are preserved.
### Internal changes
- Updated the shared sink contract to support both sync and async sink
writes.
- Added a server-side default sink backed by the existing
`EventStore`/Postgres path.
- Refactored `DirectEventIngestor` to write through `ControlEventSink`
internally.
- Preserved existing server wiring by wrapping `EventStore` inputs into
the default sink internally.
- Added test coverage proving the ingestor can accept a sink directly.
### Out of scope
- Config-driven sink selection.
- Alternate server sinks such as ClickHouse, OTEL, Kafka, or
vendor-specific sinks.
- Changes to server query/stats read-path behavior.
- Changes to SDK behavior beyond the minimal shared sink contract
compatibility needed for server support.
## Risk and Rollout
**Risk level:** Medium
**Rollback plan:**
1. Revert the shared sink contract async compatibility change.
2. Remove the new server sink adapter.
3. Restore `DirectEventIngestor` to writing directly to
`EventStore.store(...)`.
4. Keep the server endpoint and startup wiring unchanged.
## Testing
- [x] Added or updated automated tests
- [x] Ran `make check` (or explained why not)
> Validation has not been run by me in this branch flow; recommended
focused server tests should be run locally/CI.
- [ ] Manually verified behavior
## Checklist
- [ ] Linked issue/spec (if applicable)
- [x] Updated docs/examples for user-facing changes
> No docs/examples updates were needed because this story is
internal-only and preserves existing API behavior.
- [ ] Included any required follow-up tasks1 parent 07ba22f commit 6791d73
File tree
6 files changed
+90
-12
lines changed- examples/agent_control_demo
- sdks/python/src/agent_control
- server
- src/agent_control_server/observability
- ingest
- tests
- telemetry/tests
6 files changed
+90
-12
lines changed| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
154 | 154 | | |
155 | 155 | | |
156 | 156 | | |
| 157 | + | |
157 | 158 | | |
158 | 159 | | |
159 | 160 | | |
| |||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
54 | 54 | | |
55 | 55 | | |
56 | 56 | | |
57 | | - | |
| 57 | + | |
| 58 | + | |
| 59 | + | |
| 60 | + | |
| 61 | + | |
58 | 62 | | |
59 | 63 | | |
60 | 64 | | |
| |||
Lines changed: 23 additions & 11 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
1 | 1 | | |
2 | 2 | | |
3 | 3 | | |
4 | | - | |
| 4 | + | |
| 5 | + | |
| 6 | + | |
5 | 7 | | |
6 | 8 | | |
7 | 9 | | |
| |||
11 | 13 | | |
12 | 14 | | |
13 | 15 | | |
| 16 | + | |
14 | 17 | | |
| 18 | + | |
15 | 19 | | |
16 | 20 | | |
17 | 21 | | |
18 | 22 | | |
19 | 23 | | |
20 | 24 | | |
21 | 25 | | |
22 | | - | |
| 26 | + | |
23 | 27 | | |
24 | | - | |
25 | | - | |
| 28 | + | |
| 29 | + | |
26 | 30 | | |
27 | 31 | | |
28 | 32 | | |
29 | 33 | | |
30 | 34 | | |
31 | | - | |
| 35 | + | |
32 | 36 | | |
33 | 37 | | |
34 | 38 | | |
35 | | - | |
| 39 | + | |
| 40 | + | |
| 41 | + | |
| 42 | + | |
| 43 | + | |
36 | 44 | | |
37 | 45 | | |
38 | 46 | | |
39 | | - | |
| 47 | + | |
40 | 48 | | |
41 | 49 | | |
42 | | - | |
| 50 | + | |
| 51 | + | |
| 52 | + | |
| 53 | + | |
43 | 54 | | |
44 | 55 | | |
45 | 56 | | |
46 | | - | |
| 57 | + | |
47 | 58 | | |
48 | 59 | | |
49 | 60 | | |
| |||
59 | 70 | | |
60 | 71 | | |
61 | 72 | | |
62 | | - | |
63 | | - | |
| 73 | + | |
| 74 | + | |
| 75 | + | |
64 | 76 | | |
65 | 77 | | |
66 | 78 | | |
| |||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
| 1 | + | |
| 2 | + | |
| 3 | + | |
| 4 | + | |
| 5 | + | |
| 6 | + | |
| 7 | + | |
| 8 | + | |
| 9 | + | |
| 10 | + | |
| 11 | + | |
| 12 | + | |
| 13 | + | |
| 14 | + | |
| 15 | + | |
| 16 | + | |
| 17 | + | |
| 18 | + | |
| 19 | + | |
| 20 | + | |
| 21 | + | |
| 22 | + | |
| 23 | + | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
7 | 7 | | |
8 | 8 | | |
9 | 9 | | |
| 10 | + | |
10 | 11 | | |
11 | 12 | | |
12 | 13 | | |
| |||
37 | 38 | | |
38 | 39 | | |
39 | 40 | | |
| 41 | + | |
| 42 | + | |
| 43 | + | |
| 44 | + | |
| 45 | + | |
| 46 | + | |
| 47 | + | |
| 48 | + | |
| 49 | + | |
40 | 50 | | |
41 | 51 | | |
42 | 52 | | |
| |||
117 | 127 | | |
118 | 128 | | |
119 | 129 | | |
| 130 | + | |
| 131 | + | |
| 132 | + | |
| 133 | + | |
| 134 | + | |
| 135 | + | |
| 136 | + | |
| 137 | + | |
| 138 | + | |
| 139 | + | |
| 140 | + | |
| 141 | + | |
| 142 | + | |
| 143 | + | |
| 144 | + | |
| 145 | + | |
| 146 | + | |
| 147 | + | |
| 148 | + | |
| 149 | + | |
| 150 | + | |
| 151 | + | |
| 152 | + | |
| 153 | + | |
| 154 | + | |
| 155 | + | |
| 156 | + | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
5 | 5 | | |
6 | 6 | | |
7 | 7 | | |
| 8 | + | |
8 | 9 | | |
9 | 10 | | |
10 | 11 | | |
| |||
0 commit comments