nx_source_adapter.nx source
↩ module page · 400 lines · 18296 B
1// nx_source_adapter.nx -- NX-INGEST source-adapter framework ABC.
2//
3// module: nishi-core.ingest.source_adapter
4// depends: nishi-core.io.syscalls, nishi-core.io.iso8601
5// disk_kb: 8
6// capability: CORE_IO
7//
8// license_tier: PUBLIC_NISHI_SUBSTRATE
9// genealogy_id: airbyte_connector_development_kit_2024 +
10// singer_io_tap_target_specification_2017 +
11// kafka_connect_source_connector_2015 +
12// nishi_ingestion_s_class_cardinal_2026
13//
14// The SOURCE-ADAPTER FRAMEWORK ABC. Every NX-INGEST connector
15// implements this interface. Replaces Singer.io's stringly-typed
16// STREAM / SCHEMA / RECORD / STATE messages with sealed-enum
17// type-safe primitives; replaces Airbyte CDK's Python OOP with
18// NishiLang substrate composition.
19//
20// ===== Connector lifecycle ========================================
21//
22// Every adapter implements (or composes against) these primitives:
23//
24// 1. nx_<src>_describe() -- returns SourceDescriptor: kind,
25// format, license, schedule, etc.
26// 2. nx_<src>_test_connection() -- verifies upstream reachable
27// 3. nx_<src>_fetch_next_batch() -- fetch + canonicalize + emit
28// 4. nx_<src>_check_state() -- where is the cursor; what's next
29// 5. nx_<src>_handle_schema_drift() -- when upstream schema changes
30//
31// Substrate guarantees provided by the framework:
32// - bounded loops per [[feedback-bounded-loop-discipline-jpl-rule-2]]
33// - rate limit via nx_rate_limiter
34// - circuit breaker via nx_circuit_breaker
35// - dead-letter queue for malformed records
36// - lineage chain per Pillar 7 evidence tier
37// - license-tier viral propagation
38// - additive-only state + audit log per Cardinal 13
39
40// nx_safety_envelope:
41// intended_use: AUTO_APPLIED -- primitive-specific tuning queued
42// sil_target: SIL1
43// evidence: [bulk_applied_2026-05-16, see-file-comment-for-detail]
44// verdict: NOT_YET_EVALUATED
45
46import "nx_syscalls.nx"
47import "nx_iso8601.nx"
48
49// ===== SourceKind sealed enum =====================================
50//
51// What KIND of upstream are we talking to. Drives the protocol-
52// layer dispatcher. Sealed enum replaces stringly-typed "source
53// type" config strings used by Airbyte/Singer.
54
55const NX_SRC_KIND_REST_JSON: i64 = 1 // GET/POST returning JSON
56const NX_SRC_KIND_REST_XML: i64 = 2 // GET/POST returning XML
57const NX_SRC_KIND_SOAP: i64 = 3 // SOAP WSDL endpoint
58const NX_SRC_KIND_GRAPHQL: i64 = 4
59const NX_SRC_KIND_GRPC: i64 = 5
60const NX_SRC_KIND_RSS_FEED: i64 = 6
61const NX_SRC_KIND_ATOM_FEED: i64 = 7
62const NX_SRC_KIND_OAI_PMH: i64 = 8 // scholarly metadata harvesting
63const NX_SRC_KIND_BULK_CSV_DOWNLOAD: i64 = 9
64const NX_SRC_KIND_BULK_JSON_DOWNLOAD: i64 = 10
65const NX_SRC_KIND_BULK_TSV_DOWNLOAD: i64 = 11
66const NX_SRC_KIND_BULK_DARWIN_CORE: i64 = 12 // biology occurrence archives
67const NX_SRC_KIND_BULK_JATS_XML: i64 = 13 // scholarly papers
68const NX_SRC_KIND_WEB_SCRAPE_HTML: i64 = 14
69const NX_SRC_KIND_WEB_SCRAPE_JS_RENDER: i64 = 15 // needs headless browser
70const NX_SRC_KIND_SITEMAP: i64 = 16
71const NX_SRC_KIND_FTP_DIR_LIST: i64 = 17
72const NX_SRC_KIND_SFTP_DIR_LIST: i64 = 18
73const NX_SRC_KIND_SMTP_INGEST: i64 = 19 // receive email + parse
74const NX_SRC_KIND_IMAP_FETCH: i64 = 20
75const NX_SRC_KIND_WEBSOCKET_STREAM: i64 = 21
76const NX_SRC_KIND_KAFKA_TOPIC: i64 = 22 // composes sovereign event-log
77const NX_SRC_KIND_JDBC_POSTGRES: i64 = 23
78const NX_SRC_KIND_JDBC_MYSQL: i64 = 24
79const NX_SRC_KIND_JDBC_SQLITE: i64 = 25
80const NX_SRC_KIND_CDC_POSTGRES_WAL: i64 = 26 // Debezium-equivalent
81const NX_SRC_KIND_CDC_MYSQL_BINLOG: i64 = 27
82const NX_SRC_KIND_LOCAL_FILE_WATCH: i64 = 28 // inotify-equivalent
83const NX_SRC_KIND_MANUAL_TRANSCRIPTION: i64 = 29 // human-curated JSONL (PDF-only sources)
84
85func nx_src_kind_name(k: i64) -> *u8 {
86 if k == NX_SRC_KIND_REST_JSON { return "REST_JSON" }
87 if k == NX_SRC_KIND_REST_XML { return "REST_XML" }
88 if k == NX_SRC_KIND_SOAP { return "SOAP" }
89 if k == NX_SRC_KIND_GRAPHQL { return "GRAPHQL" }
90 if k == NX_SRC_KIND_GRPC { return "GRPC" }
91 if k == NX_SRC_KIND_RSS_FEED { return "RSS_FEED" }
92 if k == NX_SRC_KIND_ATOM_FEED { return "ATOM_FEED" }
93 if k == NX_SRC_KIND_OAI_PMH { return "OAI_PMH" }
94 if k == NX_SRC_KIND_BULK_CSV_DOWNLOAD { return "BULK_CSV_DOWNLOAD" }
95 if k == NX_SRC_KIND_BULK_JSON_DOWNLOAD { return "BULK_JSON_DOWNLOAD" }
96 if k == NX_SRC_KIND_BULK_TSV_DOWNLOAD { return "BULK_TSV_DOWNLOAD" }
97 if k == NX_SRC_KIND_BULK_DARWIN_CORE { return "BULK_DARWIN_CORE" }
98 if k == NX_SRC_KIND_BULK_JATS_XML { return "BULK_JATS_XML" }
99 if k == NX_SRC_KIND_WEB_SCRAPE_HTML { return "WEB_SCRAPE_HTML" }
100 if k == NX_SRC_KIND_WEB_SCRAPE_JS_RENDER { return "WEB_SCRAPE_JS_RENDER" }
101 if k == NX_SRC_KIND_SITEMAP { return "SITEMAP" }
102 if k == NX_SRC_KIND_FTP_DIR_LIST { return "FTP_DIR_LIST" }
103 if k == NX_SRC_KIND_SFTP_DIR_LIST { return "SFTP_DIR_LIST" }
104 if k == NX_SRC_KIND_SMTP_INGEST { return "SMTP_INGEST" }
105 if k == NX_SRC_KIND_IMAP_FETCH { return "IMAP_FETCH" }
106 if k == NX_SRC_KIND_WEBSOCKET_STREAM { return "WEBSOCKET_STREAM" }
107 if k == NX_SRC_KIND_KAFKA_TOPIC { return "KAFKA_TOPIC" }
108 if k == NX_SRC_KIND_JDBC_POSTGRES { return "JDBC_POSTGRES" }
109 if k == NX_SRC_KIND_JDBC_MYSQL { return "JDBC_MYSQL" }
110 if k == NX_SRC_KIND_JDBC_SQLITE { return "JDBC_SQLITE" }
111 if k == NX_SRC_KIND_CDC_POSTGRES_WAL { return "CDC_POSTGRES_WAL" }
112 if k == NX_SRC_KIND_CDC_MYSQL_BINLOG { return "CDC_MYSQL_BINLOG" }
113 if k == NX_SRC_KIND_LOCAL_FILE_WATCH { return "LOCAL_FILE_WATCH" }
114 if k == NX_SRC_KIND_MANUAL_TRANSCRIPTION { return "MANUAL_TRANSCRIPTION" }
115 return "UNKNOWN"
116}
117
118// Does this source kind require auth (API key / OAuth / etc.)?
119func nx_src_kind_requires_auth(k: i64) -> i64 {
120 if k == NX_SRC_KIND_BULK_CSV_DOWNLOAD { return 0 }
121 if k == NX_SRC_KIND_BULK_JSON_DOWNLOAD { return 0 }
122 if k == NX_SRC_KIND_SITEMAP { return 0 }
123 if k == NX_SRC_KIND_RSS_FEED { return 0 }
124 if k == NX_SRC_KIND_ATOM_FEED { return 0 }
125 if k == NX_SRC_KIND_OAI_PMH { return 0 }
126 if k == NX_SRC_KIND_WEB_SCRAPE_HTML { return 0 }
127 if k == NX_SRC_KIND_MANUAL_TRANSCRIPTION { return 0 }
128 return 1 // most everything else needs API key or similar
129}
130
131// Is this a streaming source (continuous tail) vs batch (poll-cycle)?
132func nx_src_kind_is_streaming(k: i64) -> i64 {
133 if k == NX_SRC_KIND_WEBSOCKET_STREAM { return 1 }
134 if k == NX_SRC_KIND_KAFKA_TOPIC { return 1 }
135 if k == NX_SRC_KIND_CDC_POSTGRES_WAL { return 1 }
136 if k == NX_SRC_KIND_CDC_MYSQL_BINLOG { return 1 }
137 if k == NX_SRC_KIND_LOCAL_FILE_WATCH { return 1 }
138 if k == NX_SRC_KIND_SMTP_INGEST { return 1 }
139 return 0
140}
141
142// ===== SourceFormat sealed enum (what the payload is) =============
143
144const NX_SRC_FMT_JSON: i64 = 1
145const NX_SRC_FMT_JSONL: i64 = 2
146const NX_SRC_FMT_CSV: i64 = 3
147const NX_SRC_FMT_TSV: i64 = 4
148const NX_SRC_FMT_XML: i64 = 5
149const NX_SRC_FMT_YAML: i64 = 6
150const NX_SRC_FMT_TOML: i64 = 7
151const NX_SRC_FMT_HTML: i64 = 8
152const NX_SRC_FMT_MARKDOWN: i64 = 9
153const NX_SRC_FMT_JATS_XML: i64 = 10
154const NX_SRC_FMT_DARWIN_CORE: i64 = 11
155const NX_SRC_FMT_RSS_XML: i64 = 12
156const NX_SRC_FMT_ATOM_XML: i64 = 13
157const NX_SRC_FMT_PROTOBUF: i64 = 14
158const NX_SRC_FMT_AVRO: i64 = 15
159const NX_SRC_FMT_PARQUET: i64 = 16
160const NX_SRC_FMT_MSGPACK: i64 = 17
161const NX_SRC_FMT_CBOR: i64 = 18
162const NX_SRC_FMT_PLAIN_TEXT: i64 = 19
163const NX_SRC_FMT_BINARY_FORBIDDEN: i64 = 20 // PDF/DOCX/etc; substrate REFUSES
164
165func nx_src_format_name(f: i64) -> *u8 {
166 if f == NX_SRC_FMT_JSON { return "JSON" }
167 if f == NX_SRC_FMT_JSONL { return "JSONL" }
168 if f == NX_SRC_FMT_CSV { return "CSV" }
169 if f == NX_SRC_FMT_TSV { return "TSV" }
170 if f == NX_SRC_FMT_XML { return "XML" }
171 if f == NX_SRC_FMT_YAML { return "YAML" }
172 if f == NX_SRC_FMT_TOML { return "TOML" }
173 if f == NX_SRC_FMT_HTML { return "HTML" }
174 if f == NX_SRC_FMT_MARKDOWN { return "MARKDOWN" }
175 if f == NX_SRC_FMT_JATS_XML { return "JATS_XML" }
176 if f == NX_SRC_FMT_DARWIN_CORE { return "DARWIN_CORE" }
177 if f == NX_SRC_FMT_RSS_XML { return "RSS_XML" }
178 if f == NX_SRC_FMT_ATOM_XML { return "ATOM_XML" }
179 if f == NX_SRC_FMT_PROTOBUF { return "PROTOBUF" }
180 if f == NX_SRC_FMT_AVRO { return "AVRO" }
181 if f == NX_SRC_FMT_PARQUET { return "PARQUET" }
182 if f == NX_SRC_FMT_MSGPACK { return "MSGPACK" }
183 if f == NX_SRC_FMT_CBOR { return "CBOR" }
184 if f == NX_SRC_FMT_PLAIN_TEXT { return "PLAIN_TEXT" }
185 if f == NX_SRC_FMT_BINARY_FORBIDDEN { return "BINARY_FORBIDDEN" }
186 return "UNKNOWN"
187}
188
189// Per [[feedback-no-pdfs-no-proprietary-binary-formats]]: substrate
190// refuses to ingest forbidden-binary formats.
191func nx_src_format_is_allowed(f: i64) -> i64 {
192 if f == NX_SRC_FMT_BINARY_FORBIDDEN { return 0 }
193 return 1
194}
195
196// ===== SourceDescriptor (the connector's self-description) =======
197
198struct SourceDescriptor {
199 descriptor_hk: i64,
200 source_name_ptr: *u8, // "usda_grin" / "arxiv" / etc.
201 source_version: i64, // schema version of the descriptor
202 kind: i64, // NX_SRC_KIND_*
203 format: i64, // NX_SRC_FMT_*
204 requires_auth: i64, // 0/1
205 is_streaming: i64, // 0/1
206 rate_limit_per_minute: i64, // 0 = no explicit limit
207 polite_pool_email_ptr: *u8, // for OAI-PMH style polite-pool headers
208 user_agent_ptr: *u8, // for HTTP User-Agent
209 license_tier: i64, // composite license of source
210 share_alike_required: i64,
211 attribution_required: i64,
212 schedule_days: i64, // recommended sync cadence
213 homepage_url_ptr: *u8,
214 docs_url_ptr: *u8,
215 capability_bitmask: i64, // module-registry capability tags
216 is_enabled: i64,
217}
218
219const NX_SOURCE_DESCRIPTOR_BYTES: i64 = 128 // 16 fields * 8 bytes
220
221// ===== AdapterVerdict =============================================
222
223const NX_ADAPTER_OK: i64 = 1
224const NX_ADAPTER_AUTH_FAIL: i64 = 2
225const NX_ADAPTER_NETWORK_FAIL: i64 = 3
226const NX_ADAPTER_RATE_LIMITED: i64 = 4
227const NX_ADAPTER_SCHEMA_DRIFT: i64 = 5
228const NX_ADAPTER_PARSE_FAIL: i64 = 6
229const NX_ADAPTER_FORMAT_FORBIDDEN: i64 = 7 // upstream serves PDF etc.
230const NX_ADAPTER_CIRCUIT_OPEN: i64 = 8 // breaker tripped; backing off
231const NX_ADAPTER_DEAD_LETTER_QUEUED: i64 = 9 // malformed; routed to DLQ
232const NX_ADAPTER_LICENSE_INCOMPATIBLE: i64 = 10
233const NX_ADAPTER_DONE_NO_MORE_DATA: i64 = 11 // cursor exhausted for this batch
234
235func nx_adapter_verdict_name(v: i64) -> *u8 {
236 if v == NX_ADAPTER_OK { return "OK" }
237 if v == NX_ADAPTER_AUTH_FAIL { return "AUTH_FAIL" }
238 if v == NX_ADAPTER_NETWORK_FAIL { return "NETWORK_FAIL" }
239 if v == NX_ADAPTER_RATE_LIMITED { return "RATE_LIMITED" }
240 if v == NX_ADAPTER_SCHEMA_DRIFT { return "SCHEMA_DRIFT" }
241 if v == NX_ADAPTER_PARSE_FAIL { return "PARSE_FAIL" }
242 if v == NX_ADAPTER_FORMAT_FORBIDDEN { return "FORMAT_FORBIDDEN" }
243 if v == NX_ADAPTER_CIRCUIT_OPEN { return "CIRCUIT_OPEN" }
244 if v == NX_ADAPTER_DEAD_LETTER_QUEUED { return "DEAD_LETTER_QUEUED" }
245 if v == NX_ADAPTER_LICENSE_INCOMPATIBLE { return "LICENSE_INCOMPATIBLE" }
246 if v == NX_ADAPTER_DONE_NO_MORE_DATA { return "DONE_NO_MORE_DATA" }
247 return "UNKNOWN"
248}
249
250// ===== SourceState (cursor + watermark) ===========================
251
252struct SourceState {
253 state_hk: i64,
254 descriptor_hk: i64, // FK to SourceDescriptor
255 cursor_ptr: *u8, // opaque cursor (e.g. since-timestamp / token / ID)
256 cursor_len: i64,
257 watermark_unix: i64, // upstream's "data up to" timestamp
258 last_sync_unix: i64,
259 last_sync_record_count: i64,
260 last_sync_verdict: i64, // NX_ADAPTER_*
261 consecutive_failures: i64, // backoff signal
262 total_records_to_date: i64,
263 total_bytes_to_date: i64,
264 schema_version_seen: i64, // detect upstream schema bump
265 next_scheduled_unix: i64,
266 is_paused: i64, // manual pause flag
267}
268
269const NX_SOURCE_STATE_BYTES: i64 = 104 // 13 fields * 8 bytes
270
271// ===== Constructor ================================================
272
273func nx_source_state_new(descriptor_hk: i64) -> *SourceState {
274 let raw: *u8 = sys_mmap(NX_SOURCE_STATE_BYTES)
275 let s: *SourceState = raw as *SourceState
276 s.state_hk = 0
277 s.descriptor_hk = descriptor_hk
278 s.cursor_ptr = 0 as *u8
279 s.cursor_len = 0
280 s.watermark_unix = 0
281 s.last_sync_unix = 0
282 s.last_sync_record_count = 0
283 s.last_sync_verdict = 0
284 s.consecutive_failures = 0
285 s.total_records_to_date = 0
286 s.total_bytes_to_date = 0
287 s.schema_version_seen = 0
288 s.next_scheduled_unix = 0
289 s.is_paused = 0
290 return s
291}
292
293// ===== Outcome record (per fetch_next_batch invocation) ===========
294
295struct FetchBatchOutcome {
296 outcome_hk: i64,
297 state_hk: i64,
298 batch_started_unix: i64,
299 batch_ended_unix: i64,
300 records_fetched: i64,
301 records_canonicalized: i64,
302 records_emitted: i64,
303 records_dead_lettered: i64,
304 bytes_in: i64,
305 bytes_out: i64,
306 verdict: i64, // NX_ADAPTER_*
307 new_cursor_ptr: *u8, // advance state.cursor on success
308 new_cursor_len: i64,
309 new_watermark_unix: i64,
310 schema_drift_kind: i64, // NX_SCHEMA_DRIFT_* if applicable
311}
312
313const NX_FETCH_BATCH_OUTCOME_BYTES: i64 = 112 // 14 fields * 8 bytes
314
315// ===== Apply outcome to state (post-fetch update) =================
316//
317// Composes the state machine: success advances cursor + bumps
318// total counts + clears consecutive_failures; failure bumps counter
319// (drives backoff in nx_circuit_breaker).
320
321func nx_source_state_apply_outcome(s: *SourceState, o: *FetchBatchOutcome) -> i64 {
322 if s == 0 as *SourceState { return -1 }
323 if o == 0 as *FetchBatchOutcome { return -1 }
324
325 s.last_sync_unix = o.batch_ended_unix
326 s.last_sync_record_count = o.records_emitted
327 s.last_sync_verdict = o.verdict
328 s.total_records_to_date = s.total_records_to_date + o.records_emitted
329 s.total_bytes_to_date = s.total_bytes_to_date + o.bytes_out
330
331 if o.verdict == NX_ADAPTER_OK {
332 s.consecutive_failures = 0
333 if o.new_cursor_len > 0 {
334 s.cursor_ptr = o.new_cursor_ptr
335 s.cursor_len = o.new_cursor_len
336 }
337 if o.new_watermark_unix > s.watermark_unix {
338 s.watermark_unix = o.new_watermark_unix
339 }
340 }
341 if o.verdict != NX_ADAPTER_OK {
342 if o.verdict != NX_ADAPTER_DONE_NO_MORE_DATA {
343 s.consecutive_failures = s.consecutive_failures + 1
344 }
345 }
346 return 0
347}
348
349// ===== Schema-drift kind sealed enum ==============================
350
351const NX_SCHEMA_DRIFT_NONE: i64 = 0
352const NX_SCHEMA_DRIFT_NEW_FIELD: i64 = 1 // benign; field added
353const NX_SCHEMA_DRIFT_REMOVED_FIELD: i64 = 2 // may break consumers
354const NX_SCHEMA_DRIFT_TYPE_CHANGED: i64 = 3 // breaking
355const NX_SCHEMA_DRIFT_RENAMED_FIELD: i64 = 4 // probable rename (heuristic)
356const NX_SCHEMA_DRIFT_SEMANTIC_SHIFT: i64 = 5 // type same, meaning changed
357const NX_SCHEMA_DRIFT_UNKNOWN: i64 = 6
358
359func nx_schema_drift_kind_name(k: i64) -> *u8 {
360 if k == NX_SCHEMA_DRIFT_NONE { return "NONE" }
361 if k == NX_SCHEMA_DRIFT_NEW_FIELD { return "NEW_FIELD" }
362 if k == NX_SCHEMA_DRIFT_REMOVED_FIELD { return "REMOVED_FIELD" }
363 if k == NX_SCHEMA_DRIFT_TYPE_CHANGED { return "TYPE_CHANGED" }
364 if k == NX_SCHEMA_DRIFT_RENAMED_FIELD { return "RENAMED_FIELD" }
365 if k == NX_SCHEMA_DRIFT_SEMANTIC_SHIFT { return "SEMANTIC_SHIFT" }
366 if k == NX_SCHEMA_DRIFT_UNKNOWN { return "UNKNOWN" }
367 return "INVALID"
368}
369
370// Is this drift kind breaking (caller must investigate before
371// continuing) or benign (substrate can auto-handle)?
372func nx_schema_drift_is_breaking(k: i64) -> i64 {
373 if k == NX_SCHEMA_DRIFT_NEW_FIELD { return 0 } // benign — new fields ignored
374 if k == NX_SCHEMA_DRIFT_NONE { return 0 }
375 return 1
376}
377
378// ===== Lineage record (per emitted record) ========================
379//
380// Composes [[project-wall-ingestion-tools-2026-05-15]] discipline +
381// Pillar 7 evidence-tier + license-tier propagation.
382
383struct IngestLineage {
384 lineage_hk: i64,
385 record_hk: i64, // the canonical record this lineage describes
386 source_descriptor_hk: i64,
387 source_record_id: i64, // upstream's primary key
388 fetched_at_unix: i64,
389 fetch_batch_hk: i64, // FK to FetchBatchOutcome
390 canonical_at_unix: i64,
391 emitted_at_unix: i64,
392 license_tier: i64, // per-source as ingested
393 share_alike_required: i64,
394 attribution_required: i64,
395 evidence_tier: i64, // Pillar 7
396 genealogy_id_ptr: *u8, // human-readable cite ("usda_grin_<id>")
397 is_current: i64,
398}
399
400const NX_INGEST_LINEAGE_BYTES: i64 = 104 // 13 fields * 8 bytes