nx_source_adapter.nx
buildroot/runtime/nx_source_adapter.nx
about
nx_source_adapter.nx -- NX-INGEST source-adapter framework ABC.
module: nishi-core.ingest.source_adapter
depends: nishi-core.io.syscalls, nishi-core.io.iso8601
disk_kb: 8
capability: CORE_IO
license_tier: PUBLIC_NISHI_SUBSTRATE
genealogy_id: airbyte_connector_development_kit_2024 +
singer_io_tap_target_specification_2017 +
kafka_connect_source_connector_2015 +
nishi_ingestion_s_class_cardinal_2026
The SOURCE-ADAPTER FRAMEWORK ABC. Every NX-INGEST connector
implements this interface. Replaces Singer.io's stringly-typed
STREAM / SCHEMA / RECORD / STATE messages with sealed-enum
type-safe primitives; replaces Airbyte CDK's Python OOP with
NishiLang substrate composition.
===== Connector lifecycle ========================================
Every adapter implements (or composes against) these primitives:
1. nx_<src>_describe() -- returns SourceDescriptor: kind,
format, license, schedule, etc.
2. nx_<src>_test_connection() -- verifies upstream reachable
3. nx_<src>_fetch_next_batch() -- fetch + canonicalize + emit
4. nx_<src>_check_state() -- where is the cursor; what's next
5. nx_<src>_handle_schema_drift() -- when upstream schema changes
Substrate guarantees provided by the framework:
- bounded loops per [[feedback-bounded-loop-discipline-jpl-rule-2]]
- rate limit via nx_rate_limiter
- circuit breaker via nx_circuit_breaker
- dead-letter queue for malformed records
- lineage chain per Pillar 7 evidence tier
- license-tier viral propagation
- additive-only state + audit log per Cardinal 13
dependencies 2 imports · 1 importers
imports: nx_syscalls.nxnx_iso8601.nx
imported by: nx_smoke_harness.nx
structs
| 198 | struct SourceDescriptor |
| 252 | struct SourceState |
| 295 | struct FetchBatchOutcome |
| 383 | struct IngestLineage |
consts
| 55 | const NX_SRC_KIND_REST_JSON: i64 = 1 // GET/POST returning JSON |
| 56 | const NX_SRC_KIND_REST_XML: i64 = 2 // GET/POST returning XML |
| 57 | const NX_SRC_KIND_SOAP: i64 = 3 // SOAP WSDL endpoint |
| 58 | const NX_SRC_KIND_GRAPHQL: i64 = 4 |
| 59 | const NX_SRC_KIND_GRPC: i64 = 5 |
| 60 | const NX_SRC_KIND_RSS_FEED: i64 = 6 |
| 61 | const NX_SRC_KIND_ATOM_FEED: i64 = 7 |
| 62 | const NX_SRC_KIND_OAI_PMH: i64 = 8 // scholarly metadata harvesting |
| 63 | const NX_SRC_KIND_BULK_CSV_DOWNLOAD: i64 = 9 |
| 64 | const NX_SRC_KIND_BULK_JSON_DOWNLOAD: i64 = 10 |
| 65 | const NX_SRC_KIND_BULK_TSV_DOWNLOAD: i64 = 11 |
| 66 | const NX_SRC_KIND_BULK_DARWIN_CORE: i64 = 12 // biology occurrence archives |
| 67 | const NX_SRC_KIND_BULK_JATS_XML: i64 = 13 // scholarly papers |
| 68 | const NX_SRC_KIND_WEB_SCRAPE_HTML: i64 = 14 |
| 69 | const NX_SRC_KIND_WEB_SCRAPE_JS_RENDER: i64 = 15 // needs headless browser |
| 70 | const NX_SRC_KIND_SITEMAP: i64 = 16 |
| 71 | const NX_SRC_KIND_FTP_DIR_LIST: i64 = 17 |
| 72 | const NX_SRC_KIND_SFTP_DIR_LIST: i64 = 18 |
| 73 | const NX_SRC_KIND_SMTP_INGEST: i64 = 19 // receive email + parse |
| 74 | const NX_SRC_KIND_IMAP_FETCH: i64 = 20 |
| 75 | const NX_SRC_KIND_WEBSOCKET_STREAM: i64 = 21 |
| 76 | const NX_SRC_KIND_KAFKA_TOPIC: i64 = 22 // composes sovereign event-log |
| 77 | const NX_SRC_KIND_JDBC_POSTGRES: i64 = 23 |
| 78 | const NX_SRC_KIND_JDBC_MYSQL: i64 = 24 |
| 79 | const NX_SRC_KIND_JDBC_SQLITE: i64 = 25 |
| 80 | const NX_SRC_KIND_CDC_POSTGRES_WAL: i64 = 26 // Debezium-equivalent |
| 81 | const NX_SRC_KIND_CDC_MYSQL_BINLOG: i64 = 27 |
| 82 | const NX_SRC_KIND_LOCAL_FILE_WATCH: i64 = 28 // inotify-equivalent |
| 83 | const NX_SRC_KIND_MANUAL_TRANSCRIPTION: i64 = 29 // human-curated JSONL (PDF-only sources) |
| 144 | const NX_SRC_FMT_JSON: i64 = 1 |
| 145 | const NX_SRC_FMT_JSONL: i64 = 2 |
| 146 | const NX_SRC_FMT_CSV: i64 = 3 |
| 147 | const NX_SRC_FMT_TSV: i64 = 4 |
| 148 | const NX_SRC_FMT_XML: i64 = 5 |
| 149 | const NX_SRC_FMT_YAML: i64 = 6 |
| 150 | const NX_SRC_FMT_TOML: i64 = 7 |
| 151 | const NX_SRC_FMT_HTML: i64 = 8 |
| 152 | const NX_SRC_FMT_MARKDOWN: i64 = 9 |
| 153 | const NX_SRC_FMT_JATS_XML: i64 = 10 |
| 154 | const NX_SRC_FMT_DARWIN_CORE: i64 = 11 |
| 155 | const NX_SRC_FMT_RSS_XML: i64 = 12 |
| 156 | const NX_SRC_FMT_ATOM_XML: i64 = 13 |
| 157 | const NX_SRC_FMT_PROTOBUF: i64 = 14 |
| 158 | const NX_SRC_FMT_AVRO: i64 = 15 |
| 159 | const NX_SRC_FMT_PARQUET: i64 = 16 |
| 160 | const NX_SRC_FMT_MSGPACK: i64 = 17 |
| 161 | const NX_SRC_FMT_CBOR: i64 = 18 |
| 162 | const NX_SRC_FMT_PLAIN_TEXT: i64 = 19 |
| 163 | const NX_SRC_FMT_BINARY_FORBIDDEN: i64 = 20 // PDF/DOCX/etc; substrate REFUSES |
| 219 | const NX_SOURCE_DESCRIPTOR_BYTES: i64 = 128 // 16 fields * 8 bytes |
| 223 | const NX_ADAPTER_OK: i64 = 1 |
| 224 | const NX_ADAPTER_AUTH_FAIL: i64 = 2 |
| 225 | const NX_ADAPTER_NETWORK_FAIL: i64 = 3 |
| 226 | const NX_ADAPTER_RATE_LIMITED: i64 = 4 |
| 227 | const NX_ADAPTER_SCHEMA_DRIFT: i64 = 5 |
| 228 | const NX_ADAPTER_PARSE_FAIL: i64 = 6 |
| 229 | const NX_ADAPTER_FORMAT_FORBIDDEN: i64 = 7 // upstream serves PDF etc. |
| 230 | const NX_ADAPTER_CIRCUIT_OPEN: i64 = 8 // breaker tripped; backing off |
| 231 | const NX_ADAPTER_DEAD_LETTER_QUEUED: i64 = 9 // malformed; routed to DLQ |
| 232 | const NX_ADAPTER_LICENSE_INCOMPATIBLE: i64 = 10 |
| 233 | const NX_ADAPTER_DONE_NO_MORE_DATA: i64 = 11 // cursor exhausted for this batch |
| 269 | const NX_SOURCE_STATE_BYTES: i64 = 104 // 13 fields * 8 bytes |
| 313 | const NX_FETCH_BATCH_OUTCOME_BYTES: i64 = 112 // 14 fields * 8 bytes |
| 351 | const NX_SCHEMA_DRIFT_NONE: i64 = 0 |
| 352 | const NX_SCHEMA_DRIFT_NEW_FIELD: i64 = 1 // benign; field added |
| 353 | const NX_SCHEMA_DRIFT_REMOVED_FIELD: i64 = 2 // may break consumers |
| 354 | const NX_SCHEMA_DRIFT_TYPE_CHANGED: i64 = 3 // breaking |
| 355 | const NX_SCHEMA_DRIFT_RENAMED_FIELD: i64 = 4 // probable rename (heuristic) |
| 356 | const NX_SCHEMA_DRIFT_SEMANTIC_SHIFT: i64 = 5 // type same, meaning changed |
| 357 | const NX_SCHEMA_DRIFT_UNKNOWN: i64 = 6 |
| 400 | const NX_INGEST_LINEAGE_BYTES: i64 = 104 // 13 fields * 8 bytes |
functions
| 85 | func nx_src_kind_name(k: i64) -> *u8 |
| 119 | func nx_src_kind_requires_auth(k: i64) -> i64 |
| 132 | func nx_src_kind_is_streaming(k: i64) -> i64 |
| 165 | func nx_src_format_name(f: i64) -> *u8 |
| 191 | func nx_src_format_is_allowed(f: i64) -> i64 |
| 235 | func nx_adapter_verdict_name(v: i64) -> *u8 |
| 273 | func nx_source_state_new(descriptor_hk: i64) -> *SourceState calls 1: sys_mmap |
| 321 | func nx_source_state_apply_outcome(s: *SourceState, o: *FetchBatchOutcome) -> i64 |
| 359 | func nx_schema_drift_kind_name(k: i64) -> *u8 |
| 372 | func nx_schema_drift_is_breaking(k: i64) -> i64 |