gdt-sync-worker
v1.1.1
Published
Standalone Reusable PostgreSQL-Backed GDT Tax Invoice Sync Worker Microservice
Readme
Standalone Multi-Project GDT Sync Worker Microservice
A high-performance, stateless worker microservice powered by pg-boss and PostgreSQL Pub/Sub (LISTEN/NOTIFY) for synchronizing Vietnamese tax invoices (Tổng cục Thuế / GDT) across multiple independent dashboard projects, ERPs, and accounting microservices.
🌟 Key Architecture & Highlights
┌────────────────────────────────────────────────────────────────────────┐
│ External Triggers & Crons │
│ Dashboard 1, Dashboard 2, ERP, or Scheduled Sync (Payload: MST + Pass)│
└───────────────────────────────────┬────────────────────────────────────┘
│ (Enqueue Job with deterministic singletonKey)
▼
┌────────────────────────────────────────────────────────────────────────┐
│ pg-boss Queue in PostgreSQL │
│ - Deduplication: 1 GDT job runs even if multiple dashboards trigger │
│ - Concurrency & Rate Limiting (safe delay between GDT API calls) │
│ - Automatic retries & exponential backoff │
└───────────────────────────────────┬────────────────────────────────────┘
│ (Worker Execution)
▼
┌────────────────────────────────────────────────────────────────────────┐
│ gdt-sync-worker │
│ 1. Authenticates directly with GDT API using payload (mst + pass) │
│ 2. In-memory token caching (23h) & stampede lock per MST │
│ 3. Fetches raw invoice list & XML/detail line items │
│ 4. Saves canonical raw data & attachments to PostgreSQL │
│ 5. Broadcasts real-time events over PostgreSQL pg_notify │
└───────────────────┬────────────────────────────────┬───────────────────┘
│ │
▼ (Save Canonical Record) ▼ (Broadcast pg_notify)
┌─────────────────────────────────────┐ ┌──────────────────────────────┐
│ Table: gdt_sync.invoices (JSONB) │ │ Channel: "gdt_invoice_events"│
│ - id, shdon, khhdon, khmshdon │ └──────────────┬───────────────┘
│ - seller_mst, buyer_mst, date │ │
│ - raw_data (JSONB - 139+ GDT fields)│ ┌────────┴────────┐
│ - raw_xml & xml_url │ ▼ ▼
└─────────────────────────────────────┘ ┌─────────────┐ ┌─────────────┐
│ Dashboard 1 │ │ Dashboard 2 │
│(Custom DB 1)│ │(Custom DB 2)│
└─────────────┘ └─────────────┘- Zero Data Loss Guarantee: Complete invoice headers, line items, full raw JSONB (139+ fields), and raw XML content are permanently stored in the canonical
invoicestable (gdt_invoices). Downstream listeners can go offline, crash, or restart at any time without losing any invoice records. - Automatic Database Trigger Pub/Sub (
AFTER INSERT OR UPDATE): A PostgreSQL trigger (trg_notify_invoice_event) automatically emits typed pub/sub events (INVOICE_HEADER_SYNCED,INVOICE_DETAIL_SYNCED,INVOICE_DETAIL_FAILED) overpg_notifywhenever any invoice is inserted or updated in the table (by the worker, background tasks, or direct SQL). - Deterministic Job Deduplication (
pg-boss): EnforcingsingletonKey: header_${mst}_${startDate}_${endDate}ensures that even if 5 dashboards trigger a sync for the same company at the same time, only 1 request hits GDT. - Zero Schema Coupling (Canonical Raw JSONB): Complete GDT invoice responses (139+ fields) are preserved in
raw_data JSONB. Consuming projects can access any field (e.g.ngày ký/nky,chữ ký số/nbcks) directly without modifying the sync worker. - Automatic Existing Detail Detection: Subsequent syncs detect already-synced invoices in the database and skip calling the GDT detail API, protecting against rate limits and 429 errors.
- Catch-Up & Historical Querying (
getInvoices): Reconnecting or new listeners can query missed invoices directly with filters (companyId,mst,updatedSince,startDate,endDate, etc.). - Pluggable Storage (
IStorageProvider): Built-inPgBroadcastProviderandMemoryStorageProvider.
📦 Canonical Database Schema (PostgreSQL)
The worker uses a single canonical table for invoice persistence:
CREATE TABLE IF NOT EXISTS gdt_sync.invoices (
id TEXT PRIMARY KEY, -- e.g. e86f187e-ba49-4917-863f-c2fd4c8f353b
company_id TEXT, -- Optional client-side UUID
mst TEXT, -- Taxpayer MST
shdon TEXT NOT NULL, -- Raw GDT invoice number (e.g. 8)
khhdon TEXT, -- Raw GDT invoice symbol/series (e.g. C26TYY)
khmshdon TEXT DEFAULT '1', -- Raw GDT invoice form/type (e.g. 1)
seller_mst TEXT, -- nbmst
buyer_mst TEXT, -- nmmst
date TEXT NOT NULL, -- tdlap (issue date)
total DOUBLE PRECISION DEFAULT 0, -- tgtttbso
status INTEGER DEFAULT 1, -- tthai (1: Normal, 2: Replaced, 3: Adjusted, 6: Cancelled)
detail_sync_status TEXT DEFAULT 'pending', -- pending | synced | failed
detail_error TEXT,
detail_retries INTEGER DEFAULT 0,
raw_data JSONB DEFAULT '{}'::jsonb, -- Full 139+ raw GDT fields
raw_xml TEXT, -- Dedicated raw XML invoice text column
xml_url TEXT, -- Cloudflare R2 / S3 storage URL
synced_at TIMESTAMP WITH TIME ZONE,
created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
);
CREATE INDEX IF NOT EXISTS idx_gdt_sync_inv_updated_at ON gdt_sync.invoices(updated_at);
CREATE INDEX IF NOT EXISTS idx_gdt_sync_inv_company ON gdt_sync.invoices(company_id, date);
CREATE INDEX IF NOT EXISTS idx_gdt_sync_inv_seller ON gdt_sync.invoices(seller_mst);
CREATE INDEX IF NOT EXISTS idx_gdt_sync_inv_buyer ON gdt_sync.invoices(buyer_mst);
CREATE INDEX IF NOT EXISTS idx_gdt_sync_inv_status ON gdt_sync.invoices(detail_sync_status);
CREATE INDEX IF NOT EXISTS idx_gdt_sync_inv_shdon ON gdt_sync.invoices(shdon);
CREATE INDEX IF NOT EXISTS idx_gdt_sync_inv_khhdon ON gdt_sync.invoices(khhdon);
### Automatic PostgreSQL Trigger for Real-Time Pub/Sub
The table automatically installs an `AFTER INSERT OR UPDATE` trigger function `notify_gdt_invoice_event()`. Whenever a row is inserted or updated, PostgreSQL automatically broadcasts the real-time event to the `gdt_invoice_events` channel:
```sql
-- Installed automatically by ensureTable():
CREATE TRIGGER trg_notify_invoice_event
AFTER INSERT OR UPDATE ON gdt_invoices
FOR EACH ROW
EXECUTE FUNCTION notify_gdt_invoice_event('gdt_invoice_events');🚀 How Consuming Projects Integrate
1. Triggering Invoice Sync (pg-boss Queue)
Any backend application (Next.js, NestJS, Express, FastAPI, Laravel) enqueues a sync job into PostgreSQL without needing direct HTTP access to the worker:
import { PgBoss } from "pg-boss";
const boss = new PgBoss({
connectionString: process.env.DATABASE_URL || "postgres://postgres:postgres@localhost:5432/hddtv2",
schema: process.env.PGBOSS_QUEUE_SCHEMA || process.env.PGBOSS_SCHEMA || "gdt_sync",
});
await boss.start();
// Trigger a sync for a company:
await boss.send("gdt-header-sync", {
mst: "0108112848",
password: "Tcvd@2026",
companyId: "my_dashboard_company_uuid", // Optional
startDate: "2026-07-01",
endDate: "2026-07-01",
invoiceTypes: ["purchase", "sold"],
queryTypes: ["query"],
}, {
// Deterministic singletonKey prevents duplicate syncs across multiple dashboards
singletonKey: `header_0108112848_2026-07-01_2026-07-01`,
});2. Real-Time Event Subscription (GdtEventListener)
Consuming projects subscribe to real-time events over PostgreSQL LISTEN / NOTIFY to execute custom domain logic (e.g. customer deduplication, product catalog linking, warehouse stock updates):
import { GdtEventListener, type GdtInvoiceEvent } from "gdt-sync-worker";
const listener = new GdtEventListener(process.env.DATABASE_URL);
listener.on(async (event: GdtInvoiceEvent) => {
if (event.event === "INVOICE_DETAIL_SYNCED") {
console.log(`⚡ Received invoice ${event.invoiceNumber} for MST ${event.mst}`);
const rawInv = event.rawData;
const isBuy = event.invoiceType === "purchase";
const partnerMst = isBuy ? rawInv.nbmst : rawInv.nmmst;
// --- Dashboard Custom Business Logic ---
// 1. Detect if customer / supplier already exists in Dashboard's custom table:
let customer = await myDb.query(
"SELECT id FROM my_dashboard_customers WHERE tax_code = $1",
[partnerMst]
);
let customerId = customer.rows[0]?.id;
if (!customerId) {
const created = await myDb.query(
"INSERT INTO my_dashboard_customers (name, tax_code, phone) VALUES ($1, $2, $3) RETURNING id",
[rawInv.nmten || "Khách hàng mới", partnerMst, rawInv.nmsdthoai]
);
customerId = created.rows[0].id;
}
// 2. Insert into Dashboard's custom invoices table:
await myDb.query(
`INSERT INTO my_dashboard_invoices (customer_id, total, ngay_ky, raw_payload)
VALUES ($1, $2, $3, $4)`,
[customerId, event.total, rawInv.nky, JSON.stringify(rawInv)]
);
}
});
// Start listening
await listener.start();3. Querying Invoices Directly from PostgreSQL
Because the worker stores all raw data in raw_data JSONB, any project can query fields dynamically:
// Query in your Dashboard backend:
const { rows } = await pool.query(`
SELECT
shdon,
khhdon,
khmshdon,
total,
raw_data->>'nky' AS ngay_ky,
raw_data->>'nbcks' AS chu_ky_nguoi_ban,
raw_data->'hdhhdvu' AS danh_sach_san_pham,
raw_xml
FROM gdt_sync.invoices
WHERE seller_mst = $1 OR buyer_mst = $1
`, ["0108112848"]);4. Catching Up on Missed Invoices After Listener Downtime
If your listener was offline, restarted, or newly deployed, use storage.getInvoices() or direct SQL to fetch all invoices synced during the downtime:
import { PgBroadcastProvider, pool } from "gdt-sync-worker";
const storage = new PgBroadcastProvider(pool);
// Fetch all invoices updated in the last 2 hours while listener was down:
const missedInvoices = await storage.getInvoices({
companyId: "my_dashboard_company_uuid",
updatedSince: new Date(Date.now() - 2 * 60 * 60 * 1000),
detailSyncStatus: "synced",
});
console.log(`Retrieved ${missedInvoices.length} missed invoices!`);🛠️ Running the Sync Worker
1. Local Development
# Install dependencies
pnpm install
# Run unit tests
pnpm exec tsx test/unit-test.ts
# Run pub/sub real-time event tests
pnpm exec tsx test/events-test.ts
# Run live GDT API test (0108112848 / Tcvd@2026)
pnpm exec tsx test/test-real-gdt.ts
# Run full live E2E test (save + notify + duplicate skip verification)
pnpm exec tsx test/e2e-real-gdt.ts
# Run PostgreSQL integration suite
pnpm test2. Production Start
# Build TypeScript
pnpm build
# Start worker service
DATABASE_URL="postgres://user:pass@host:5432/dbname" pnpm start
# Start worker with @pg-boss/dashboard web UI (port 3000)
ENABLE_DASHBOARD=true pnpm start⚙️ Environment Variables
| Variable | Default | Description |
|---|---|---|
| DATABASE_URL | postgres://postgres:postgres@localhost:5432/hddtv2 | PostgreSQL connection string |
| PGBOSS_QUEUE_SCHEMA | gdt_sync | Schema for pg-boss internal tables |
| PGBOSS_SCHEMA | gdt_sync | Alias schema for pg-boss & web dashboard |
| GDT_INVOICE_TABLE | gdt_invoices | Target table for raw invoice persistence |
| GDT_EVENTS_CHANNEL | gdt_invoice_events | PostgreSQL LISTEN/NOTIFY pub/sub channel |
| PROXY_URL | "" | SOCKS5 proxy URL for GDT API requests |
| ENABLE_DASHBOARD | false | Enable web dashboard on port 3000 |
| PGBOSS_DASHBOARD_PORT | 3000 | Port for web dashboard UI |
| R2_ACCESS_KEY_ID | "" | R2 / S3 Access Key ID (Optional) |
| R2_SECRET_ACCESS_KEY | "" | R2 / S3 Secret Access Key (Optional) |
| R2_BUCKET_NAME | hddt-xml-invoices | R2 Bucket Name (Optional) |
📋 Types Reference (src/types.ts)
InvoiceType:"purchase" | "sold"InvoiceQueryType:"query" | "sco-query"HeaderSyncJobPayload:{ mst: string; password: string; companyId?: string; startDate?: string; endDate?: string; invoiceTypes?: InvoiceType[]; queryTypes?: InvoiceQueryType[] }InvoiceDetailSyncJobPayload:{ mst: string; password?: string; companyId?: string; invoiceId: string; khhdon: string; shdon: string | number; khmshdon?: number | string; nbmst: string; invoiceType: InvoiceType; queryType?: InvoiceQueryType; rawHeader?: GdtRawInvoice }GdtRawInvoice: Full 139+ fields interface representing raw GDT header and line-item response.GdtInvoiceEvent: Typed union ofINVOICE_HEADER_SYNCED,INVOICE_DETAIL_SYNCED, andINVOICE_DETAIL_FAILED.
