Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
434 changes: 335 additions & 99 deletions services/go/asean-sw-service/cmd/main.go

Large diffs are not rendered by default.

4 changes: 4 additions & 0 deletions services/go/asean-sw-service/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ require (
github.com/IBM/sarama v1.60.2
github.com/gin-gonic/gin v1.10.0
github.com/google/uuid v1.6.0
github.com/jackc/pgx/v5 v5.9.2
github.com/prometheus/client_golang v1.24.1
github.com/sony/gobreaker v1.0.0
go.opentelemetry.io/otel v1.46.0
Expand Down Expand Up @@ -39,6 +40,9 @@ require (
github.com/grpc-ecosystem/go-grpc-middleware/v2 v2.3.2 // indirect
github.com/grpc-ecosystem/grpc-gateway/v2 v2.30.0 // indirect
github.com/hashicorp/go-uuid v1.0.3 // indirect
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
github.com/jackc/puddle/v2 v2.2.2 // indirect
github.com/jcmturner/aescts/v2 v2.0.0 // indirect
github.com/jcmturner/dnsutils/v2 v2.0.0 // indirect
github.com/jcmturner/gofork v1.7.6 // indirect
Expand Down
8 changes: 8 additions & 0 deletions services/go/asean-sw-service/go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,14 @@ github.com/grpc-ecosystem/grpc-gateway/v2 v2.30.0/go.mod h1:zOBXOsUaBSjKgmH4OGzV
github.com/hashicorp/go-uuid v1.0.2/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro=
github.com/hashicorp/go-uuid v1.0.3 h1:2gKiV6YVmrJ1i2CKKa9obLvRieoRGviZFL26PcT/Co8=
github.com/hashicorp/go-uuid v1.0.3/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro=
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM=
github.com/jackc/pgx/v5 v5.9.2 h1:3ZhOzMWnR4yJ+RW1XImIPsD1aNSz4T4fyP7zlQb56hw=
github.com/jackc/pgx/v5 v5.9.2/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4=
github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo=
github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4=
github.com/jcmturner/aescts/v2 v2.0.0 h1:9YKLH6ey7H4eDBXW8khjYslgyqG2xZikXP0EQFKrle8=
github.com/jcmturner/aescts/v2 v2.0.0/go.mod h1:AiaICIRyfYg35RUkr8yESTqvSy7csK90qZ5xfvvsoNs=
github.com/jcmturner/dnsutils/v2 v2.0.0 h1:lltnkeZGL0wILNvrNiVCR6Ro5PGU/SeBvVO/8c/iPbo=
Expand Down
179 changes: 161 additions & 18 deletions services/go/audit-service/cmd/main.go
Original file line number Diff line number Diff line change
@@ -1,9 +1,17 @@
// audit-service — Post-Clearance Audit microservice
// Implements WCO-aligned post-clearance audit with risk-weighted random
// selection, duty discrepancy calculation, and penalty notice generation.
//
// Phase 26 F3: audit state is persisted to PostgreSQL via pgx/v5 against the
// existing post_clearance_audits table (drizzle schema, columns reused
// as-is). DATABASE_URL is REQUIRED — the service fails closed at startup
// without it. Selected audits are inserted as real rows; discrepancy and
// penalty endpoints update existing audit rows and 404 honestly when no
// audit exists. No simulated data remains.
package main

import (
"context"
"fmt"
"log"
"math"
Expand All @@ -16,6 +24,7 @@ import (

"github.com/gin-gonic/gin"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgxpool"
)

// ─── Domain types ─────────────────────────────────────────────────────────────
Expand All @@ -27,8 +36,8 @@ type RiskProfile struct {
HSCode string `json:"hs_code"`
DeclaredValue float64 `json:"declared_value"`
DutyPaid float64 `json:"duty_paid"`
RiskScore float64 `json:"risk_score"` // 0–100
RiskLane string `json:"risk_lane"` // green_lane | yellow_lane | red_lane
RiskScore float64 `json:"risk_score"` // 0–100
RiskLane string `json:"risk_lane"` // green_lane | yellow_lane | red_lane
PriorAudits int `json:"prior_audits"`
PriorViolations int `json:"prior_violations"`
}
Expand All @@ -46,18 +55,18 @@ type DiscrepancyInput struct {
UCR string `json:"ucr"`
DeclaredValue float64 `json:"declared_value"`
AuditedValue float64 `json:"audited_value"`
DutyRate float64 `json:"duty_rate"` // e.g. 0.20 for 20%
DutyRate float64 `json:"duty_rate"` // e.g. 0.20 for 20%
DutyPaid float64 `json:"duty_paid"`
}

type DiscrepancyResult struct {
ValueDifference float64 `json:"value_difference"`
ValueDiffPct float64 `json:"value_diff_pct"`
AdditionalDutyOwed float64 `json:"additional_duty_owed"`
PenaltyMultiplier float64 `json:"penalty_multiplier"`
PenaltyAmount float64 `json:"penalty_amount"`
Outcome string `json:"outcome"` // compliant | minor_discrepancy | major_discrepancy | fraud_suspected
FindingSummary string `json:"finding_summary"`
ValueDifference float64 `json:"value_difference"`
ValueDiffPct float64 `json:"value_diff_pct"`
AdditionalDutyOwed float64 `json:"additional_duty_owed"`
PenaltyMultiplier float64 `json:"penalty_multiplier"`
PenaltyAmount float64 `json:"penalty_amount"`
Outcome string `json:"outcome"` // compliant | minor_discrepancy | major_discrepancy | fraud_suspected
FindingSummary string `json:"finding_summary"`
}

type PenaltyNotice struct {
Expand All @@ -73,6 +82,8 @@ type PenaltyNotice struct {
Narrative string `json:"narrative"`
}

var db *pgxpool.Pool

// ─── Selection algorithm ──────────────────────────────────────────────────────

// selectForAudit implements WCO risk-weighted random selection.
Expand Down Expand Up @@ -162,6 +173,19 @@ func selectForAudit(p RiskProfile) SelectionResult {
}
}

// recordSelectedAudit persists a newly selected audit into
// post_clearance_audits (audit_number, declaration_id, declaration_number,
// trader_id, status, outcome, trigger_reason, declared_value).
func recordSelectedAudit(ctx context.Context, p RiskProfile, res SelectionResult) error {
_, err := db.Exec(ctx, `
INSERT INTO post_clearance_audits
(audit_number, declaration_id, declaration_number, trader_id,
status, outcome, trigger_reason, declared_value)
VALUES ($1, $2, $3, $4, 'scheduled', 'pending', $5, $6)`,
res.AuditID, p.DeclarationID, p.UCR, p.TraderID, res.TriggerReason, p.DeclaredValue)
return err
}

// ─── Discrepancy calculator ───────────────────────────────────────────────────

func calcDiscrepancy(inp DiscrepancyInput) DiscrepancyResult {
Expand Down Expand Up @@ -209,6 +233,34 @@ func calcDiscrepancy(inp DiscrepancyInput) DiscrepancyResult {
}
}

// applyDiscrepancy records the audit outcome onto the latest
// post_clearance_audits row for the declaration. Returns false when no audit
// exists — the handler then fails closed with 404 rather than inventing one.
func applyDiscrepancy(ctx context.Context, inp DiscrepancyInput, res DiscrepancyResult) (bool, error) {
tag, err := db.Exec(ctx, `
UPDATE post_clearance_audits
SET audited_value = $1,
value_difference = $2,
additional_duty_assessed = $3,
penalty_amount = $4,
outcome = $5,
findings = $6,
status = 'completed',
completed_at = now(),
updated_at = now()
WHERE id = (
SELECT id FROM post_clearance_audits
WHERE declaration_id = $7
ORDER BY created_at DESC LIMIT 1
)`,
inp.AuditedValue, res.ValueDifference, res.AdditionalDutyOwed,
res.PenaltyAmount, res.Outcome, res.FindingSummary, inp.DeclarationID)
if err != nil {
return false, err
}
return tag.RowsAffected() > 0, nil
}

// ─── Penalty notice generator ─────────────────────────────────────────────────

func issuePenaltyNotice(auditID, ucr string, traderID int, disc DiscrepancyResult) PenaltyNotice {
Expand Down Expand Up @@ -239,9 +291,31 @@ func issuePenaltyNotice(auditID, ucr string, traderID int, disc DiscrepancyResul
}
}

// applyPenaltyNotice records the issued penalty against the persisted audit.
// Returns false when the audit does not exist.
func applyPenaltyNotice(ctx context.Context, n PenaltyNotice) (bool, error) {
tag, err := db.Exec(ctx, `
UPDATE post_clearance_audits
SET penalty_amount = $1,
officer_notes = $2,
status = 'escalated',
updated_at = now()
WHERE audit_number = $3`,
n.PenaltyAmount, n.Narrative, n.AuditID)
if err != nil {
return false, err
}
return tag.RowsAffected() > 0, nil
}

// ─── HTTP handlers ────────────────────────────────────────────────────────────

func handleHealth(c *gin.Context) {
var one int
if err := db.QueryRow(c.Request.Context(), "SELECT 1").Scan(&one); err != nil {
c.JSON(http.StatusServiceUnavailable, gin.H{"status": "unavailable", "service": "audit-service", "error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"status": "ok", "service": "audit-service", "ts": time.Now().UTC()})
}

Expand All @@ -252,6 +326,12 @@ func handleSelect(c *gin.Context) {
return
}
result := selectForAudit(p)
if result.Selected {
if err := recordSelectedAudit(c.Request.Context(), p, result); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "failed to persist selected audit: " + err.Error()})
return
}
}
c.JSON(http.StatusOK, result)
}

Expand All @@ -262,14 +342,23 @@ func handleDiscrepancy(c *gin.Context) {
return
}
result := calcDiscrepancy(inp)
updated, err := applyDiscrepancy(c.Request.Context(), inp, result)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "failed to persist discrepancy: " + err.Error()})
return
}
if !updated {
c.JSON(http.StatusNotFound, gin.H{"error": fmt.Sprintf("no post-clearance audit exists for declaration %d; run /api/audit/select first", inp.DeclarationID)})
return
}
c.JSON(http.StatusOK, result)
}

func handlePenalty(c *gin.Context) {
var req struct {
AuditID string `json:"audit_id"`
UCR string `json:"ucr"`
TraderID int `json:"trader_id"`
AuditID string `json:"audit_id"`
UCR string `json:"ucr"`
TraderID int `json:"trader_id"`
Disc DiscrepancyResult `json:"discrepancy"`
}
if err := c.ShouldBindJSON(&req); err != nil {
Expand All @@ -281,6 +370,15 @@ func handlePenalty(c *gin.Context) {
return
}
notice := issuePenaltyNotice(req.AuditID, req.UCR, req.TraderID, req.Disc)
updated, err := applyPenaltyNotice(c.Request.Context(), notice)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "failed to persist penalty notice: " + err.Error()})
return
}
if !updated {
c.JSON(http.StatusNotFound, gin.H{"error": "audit not found: " + req.AuditID})
return
}
c.JSON(http.StatusCreated, notice)
}

Expand All @@ -293,6 +391,12 @@ func handleBatchSelect(c *gin.Context) {
results := make([]SelectionResult, len(profiles))
for i, p := range profiles {
results[i] = selectForAudit(p)
if results[i].Selected {
if err := recordSelectedAudit(c.Request.Context(), p, results[i]); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": fmt.Sprintf("failed to persist audit for declaration %d: %v", p.DeclarationID, err)})
return
}
}
}
selected := 0
for _, r := range results {
Expand All @@ -308,25 +412,64 @@ func handleBatchSelect(c *gin.Context) {
}

func handleSelectionStats(c *gin.Context) {
// Return selection rate statistics by risk lane for dashboard display
// Selection policy metadata (static configuration) plus real audit counts
// from post_clearance_audits.
byStatus := map[string]int{}
rows, err := db.Query(c.Request.Context(), `SELECT status, COUNT(*) FROM post_clearance_audits GROUP BY status`)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
defer rows.Close()
total := 0
for rows.Next() {
var s string
var n int
if err := rows.Scan(&s, &n); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
byStatus[s] = n
total += n
}
c.JSON(http.StatusOK, gin.H{
"selection_rates": gin.H{
"green_lane": "5% base + risk modifiers",
"yellow_lane": "20% base + risk modifiers",
"red_lane": "100% (mandatory)",
},
"penalty_matrix": []gin.H{
{"range": "< 5% discrepancy", "outcome": "compliant", "multiplier": 0},
{"range": "5–20% discrepancy", "outcome": "minor_discrepancy", "multiplier": 1},
{"range": "20–50% discrepancy","outcome": "major_discrepancy", "multiplier": 2},
{"range": "> 50% discrepancy", "outcome": "fraud_suspected", "multiplier": 4},
{"range": "< 5% discrepancy", "outcome": "compliant", "multiplier": 0},
{"range": "5–20% discrepancy", "outcome": "minor_discrepancy", "multiplier": 1},
{"range": "20–50% discrepancy", "outcome": "major_discrepancy", "multiplier": 2},
{"range": "> 50% discrepancy", "outcome": "fraud_suspected", "multiplier": 4},
},
"audits_total": total,
"audits_by_status": byStatus,
})
}

// ─── Main ─────────────────────────────────────────────────────────────────────

func main() {
// Phase 26 F3: DATABASE_URL is mandatory. This service persists audits to
// post_clearance_audits and refuses to start without a database.
dsn := os.Getenv("DATABASE_URL")
if dsn == "" {
log.Fatal("DATABASE_URL is required; refusing to start without PostgreSQL persistence")
}
var err error
db, err = pgxpool.New(context.Background(), dsn)
if err != nil {
log.Fatalf("DATABASE_URL parse failed: %v", err)
}
defer db.Close()
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
if err := db.Ping(ctx); err != nil {
log.Fatalf("PostgreSQL unreachable at startup (fail-closed): %v", err)
}

port := os.Getenv("PORT")
if port == "" {
port = "8094"
Expand Down
4 changes: 4 additions & 0 deletions services/go/audit-service/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ require (
github.com/IBM/sarama v1.60.2
github.com/gin-gonic/gin v1.10.0
github.com/google/uuid v1.6.0
github.com/jackc/pgx/v5 v5.9.2
github.com/prometheus/client_golang v1.24.1
github.com/sony/gobreaker v1.0.0
go.opentelemetry.io/otel v1.46.0
Expand Down Expand Up @@ -39,6 +40,9 @@ require (
github.com/grpc-ecosystem/go-grpc-middleware/v2 v2.3.2 // indirect
github.com/grpc-ecosystem/grpc-gateway/v2 v2.30.0 // indirect
github.com/hashicorp/go-uuid v1.0.3 // indirect
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
github.com/jackc/puddle/v2 v2.2.2 // indirect
github.com/jcmturner/aescts/v2 v2.0.0 // indirect
github.com/jcmturner/dnsutils/v2 v2.0.0 // indirect
github.com/jcmturner/gofork v1.7.6 // indirect
Expand Down
8 changes: 8 additions & 0 deletions services/go/audit-service/go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,14 @@ github.com/grpc-ecosystem/grpc-gateway/v2 v2.30.0/go.mod h1:zOBXOsUaBSjKgmH4OGzV
github.com/hashicorp/go-uuid v1.0.2/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro=
github.com/hashicorp/go-uuid v1.0.3 h1:2gKiV6YVmrJ1i2CKKa9obLvRieoRGviZFL26PcT/Co8=
github.com/hashicorp/go-uuid v1.0.3/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro=
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM=
github.com/jackc/pgx/v5 v5.9.2 h1:3ZhOzMWnR4yJ+RW1XImIPsD1aNSz4T4fyP7zlQb56hw=
github.com/jackc/pgx/v5 v5.9.2/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4=
github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo=
github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4=
github.com/jcmturner/aescts/v2 v2.0.0 h1:9YKLH6ey7H4eDBXW8khjYslgyqG2xZikXP0EQFKrle8=
github.com/jcmturner/aescts/v2 v2.0.0/go.mod h1:AiaICIRyfYg35RUkr8yESTqvSy7csK90qZ5xfvvsoNs=
github.com/jcmturner/dnsutils/v2 v2.0.0 h1:lltnkeZGL0wILNvrNiVCR6Ro5PGU/SeBvVO/8c/iPbo=
Expand Down
Loading
Loading