Enterprise services and data collection
FlowPort 0.1.6 adds 10 native Go Sources, five Destinations and Grafana Cloud / New Relic OTLP presets: 51 Source types, 39 Destination types, two input and ten output presets. Installation and runtime require no Node.js or Vector.
Azure Service Bus
Use azure_service_bus with servicebus_namespace and Worker Entra identity or connection_string. Sources select queue or topic + subscription; Destinations select one queue or topic. Pre-create entities. Receive one message with PeekLock, renew every five seconds, and Complete only after local durable admission. Capacity rejection, lease loss and cancellation do not acknowledge. Publish negotiated batches and wait for receipts. Message IDs derive from durable batch identity and row content; duplicate detection works only when enabled on the entity. Partial batch success followed by queue retry can duplicate records. Diagnostics use Peek or an empty sender batch and never consume or publish business messages. Sampling observes running input.
SQL queries and incremental collection
postgresql_query, mysql_query and sqlserver_query use a credential-free url, separate username/password, query and interval_seconds (default 30). batch_size is 1–10000, default 1000. Accept one SELECT, no SELECT INTO or multiple statements. PostgreSQL/MySQL use read-only transactions; SQL Server requires a SELECT-only account. Without cursor_column every poll repeats the full query. Incremental queries bind {{cursor}} and {{limit}}, order a unique cursor ascending, and persist each page atomically with its next cursor. Equal timestamps require tie_breaker_column, tie_breaker_start and {{tie_breaker}} with composite ordering. Query/resource/cursor changes create a new progress namespace. Polling cannot capture deletes or updates without a changed watermark; use CDC for those. SQL Server uses SELECT TOP ({{limit}}).
SELECT id, message FROM events
WHERE id > {{cursor}} ORDER BY id LIMIT {{limit}}PostgreSQL CDC
postgresql_cdc requires url, an existing pgoutput slot, publication, replication permission and wal_level=logical. Use REPLICA IDENTITY FULL for complete old update/delete values; default identity supplies available keys. unchanged_toast marks unavailable unchanged TOAST values and is not null. Durable progress overrides start_lsn or the slot position. Buffer insert/update/delete/truncate through commit, then persist transaction and end LSN atomically. Relation messages refresh column mappings; no downstream DDL runs. Confirm only the transferable position after all destinations finish. Monitor retained WAL and lag.
initial_snapshot uses a temporary logical slot and exported consistent snapshot, honoring publication columns, row filters and partition scope, then starts streaming at the snapshot LSN. This snapshot path requires PostgreSQL 15+. User-owned persistent slots are never automatically created or deleted.
MySQL CDC
mysql_cdc requires MySQL 8.0+, a unique server_id, replication permission, ROW binlog format, FULL row image and FULL row metadata. tables contains database.table per line; blank collects all row events. Resume file/position or gtid_set from durable progress; when no start is supplied, fix and persist the current binlog position. Accept transactions and progress only at commit. Emit relevant table DDL without executing it downstream. Column metadata refreshes subsequent rows. Statement-mode data and partial JSON update encoding fail explicitly. Retain sufficient binlog; purged history is not repaired automatically.
initial_snapshot requires explicit InnoDB tables. A short global read lock fixes the consistent view and binlog position, then releases before bulk reading. It requires FLUSH_TABLES/RELOAD permissions and can briefly block writes. Both CDC sources default to a 32 MiB transaction buffer (1 KiB–128 MiB) and 500-row snapshot pages. Oversized transactions stop without acknowledgement or silent splitting. Events contain operation, schema/table, before/after and protocol progress. Snapshots emit snapshot_begin, snapshot, snapshot_complete and snapshot_id. An interrupted snapshot restarts with a new ID and may repeat accepted rows; materialized consumers must reconcile generation and completion boundaries. No end-to-end exactly-once claim. Sampling never opens another change stream or creates a slot.
{"type":"postgresql_cdc","config":{"url":"postgres://HOST:5432/DB","username":"replication_reader","password":"REPLACE_ME","slot":"flowport_slot","publication":"flowport_pub","initial_snapshot":true,"max_transaction_bytes":33554432}}Microsoft 365 and Google Workspace audits
microsoft365_audit needs tenant_id and client_id/client_secret, or a short-lived token. content_type defaults to Audit.General. Enable unified auditing and subscriptions first. Only create_subscription=true creates a subscription during collection; diagnostics do not. google_workspace_audit uses credentials_json with domain-wide delegation, delegated_user with audit permission, application_name (default login), user_key (default all), and admin.reports.audit.readonly scope. Service account tokens refresh; manually supplied short-lived tokens do not.
Poll every 60 seconds, initially looking back one hour. Optional start_time is RFC3339. Windows are at most 24 hours with a two-minute delay; lookback_seconds is 1–86400. Deduplicate content/event IDs, at most 100000 retained IDs and 100 pages per pass. Persist pages and progress together; resume after restart. Content and next-page URLs must share the API origin. Provider history, lateness, licensing, permissions and quotas apply; events beyond the lookback need backfill. Sampling observes running input.
SNMP and SNMP Trap
snmp uses remote address=host:161 and numeric oids, one per line. Support v1, v2c/community and v3 authPriv with SHA256/AES128, username, auth_password and privacy_password. GET each OID or walk=true; cap variables at 1000 by default (1–10000). Preserve OID, host and value; nonnumeric variables are not numeric metrics. snmp_trap listens on UDP, default 127.0.0.1:1162. Versions and credentials must match; v3 also needs the appropriate hexadecimal engine_id. Traps have no reliable acknowledgement or replay. Inform replies follow durable admission only. API/UI mask credentials and event fields exclude them.
Azure Monitor Logs
azure_monitor_logs needs an HTTPS DCR/DCE base url, immutable dcr_id and stream_name. Pre-create the table and DCR and grant publishing permission. Use Worker Entra identity, tenant_id/client_id/client_secret or a short-lived Token. Send JSON arrays capped at 512 KiB. Default columns are measurement/tags/fields/timestamp/category; the DCR declaration must match. Advanced field_mapping maps columns from fields.*, tags.*, measurement, timestamp or category, followed by DCR transformations. Identity probes never ingest and do not prove DCR write permission.
BigQuery Storage Write
gcp_bigquery needs project, dataset and table with Worker ADC, credentials_json or a short-lived Token. Use Storage Write API _default, real table schema and protobuf, checking each AppendRows receipt with requests capped at 4 MiB. Default schema: measurement STRING, tags JSON, fields JSON, timestamp INT64 nanoseconds, category STRING. field_mapping selects real columns. Support STRING, JSON, INT64, DOUBLE, BOOL, BYTES and TIMESTAMP (microseconds); complex/repeated columns fail explicitly. JSON preserves large-integer text. The default stream is at-least-once, without offset-based retry deduplication. Use stable business IDs downstream. Row errors are failures; no automatic table creation. Diagnostics read schema only.
AWS Data Firehose
aws_firehose requires region and stream_name, with Worker credentials or access_key/secret_key/session_token. This differs from Firehose HTTP input and Kinesis Data Streams. PutRecordBatch sends at most 500 records/4 MiB, at most 1000 KiB per record. Check each receipt and retry only transient failures, up to five rounds. Queue retries may repeat successes. Acceptance by Firehose does not prove final downstream delivery. DescribeDeliveryStream diagnostics publish no records.
Snowflake Streaming
snowflake needs HTTPS url, database, schema and table. token_type=scoped uses the ingest base URL and existing Scoped Token; oauth uses the account URL and existing OAuth Token; jwt uses account, username and an RSA private_key PEM of at least 2048 bits. Go signs JWT, discovers the ingest host and exchanges a Scoped Token. Explicit trusted HTTPS ingest_url supports private connectivity with certificate verification. Configuration encrypts and API-masks tokens and keys.
Use Elastic Channel REST NDJSON with raw batches capped at 3 MiB, optional gzip and field_mapping. Stable request IDs derive from durable batch identity and content. Require message=OK after HTTP success. This confirms durable buffering, not immediate table visibility or successful row transformations; inspect Snowflake error tables. Elastic is at-least-once with no ordering guarantee. Ambiguous retries can duplicate rows; no Named Channel exactly-once claim. Scoped mode has no safe write probe, so diagnostics preserve an unknown scope.
Grafana Cloud and New Relic presets
Reuse OTLP HTTP for logging/metric/tracing with gzip. Grafana Cloud uses its console OTLP base URL including /otlp, Basic username and token password. New Relic uses an Ingest License Key in API Key, sent as api-key, and the regional endpoint. Presets do not certify account interoperability.
Validation scope
Real PostgreSQL and MySQL on OrbStack cover incremental queries, consistent initial snapshots, inserts/updates/deletes and restart. Protocol tests cover audit admission/deduplication, Inform acknowledgement, Firehose partial failures, BigQuery protobuf/rejection, Azure DCR requests and Snowflake receipts. SQL Server and real cloud IAM, quotas, table schemas, network, service versions and sustained load need deployment-specific tests. Connection diagnostics and Pipeline samples do not replace end-to-end delivery.