npm package discovery and stats viewer.

Discover Tips

  • General search

    [free text search, go nuts!]

  • Package details

    pkg:[package-name]

  • User packages

    @[username]

Sponsor

Optimize Toolset

I’ve always been into building performant and accessible sites, but lately I’ve been taking it extremely seriously. So much so that I’ve been building a tool to help me optimize and monitor the sites that I build to make sure that I’m making an attempt to offer the best experience to those who visit them. If you’re into performant, accessible and SEO friendly sites, you might like it too! You can check it out at Optimize Toolset.

About

Hi, 👋, I’m Ryan Hefner  and I built this site for me, and you! The goal of this site was to provide an easy way for me to check the stats on my npm packages, both for prioritizing issues and updates, and to give me a little kick in the pants to keep up on stuff.

As I was building it, I realized that I was actually using the tool to build the tool, and figured I might as well put this out there and hopefully others will find it to be a fast and useful way to search and browse npm packages as I have.

If you’re interested in other things I’m working on, follow me on Twitter or check out the open source projects I’ve been publishing on GitHub.

I am also working on a Twitter bot for this site to tweet the most popular, newest, random packages from npm. Please follow that account now and it will start sending out packages soon–ish.

Open Software & Tools

This site wouldn’t be possible without the immense generosity and tireless efforts from the people who make contributions to the world and share their work via open source initiatives. Thank you 🙏

© 2026 – Pkg Stats / Ryan Hefner

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 invoices table (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) over pg_notify whenever any invoice is inserted or updated in the table (by the worker, background tasks, or direct SQL).
  • Deterministic Job Deduplication (pg-boss): Enforcing singletonKey: 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-in PgBroadcastProvider and MemoryStorageProvider.

📦 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 test

2. 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 of INVOICE_HEADER_SYNCED, INVOICE_DETAIL_SYNCED, and INVOICE_DETAIL_FAILED.