code wiki / (root) / nx_source_adapter.nx

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