database

package module
v0.0.0-...-363cbd8 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Nov 12, 2025 License: AGPL-3.0 Imports: 13 Imported by: 0

README

DictaMesh Database Package

Comprehensive database infrastructure for DictaMesh framework with advanced features for data source agnostic metadata management, vector search, caching, and compliance.

Features

Core Database Management
  • Multiple Connection Pools: pgx for performance, GORM for ORM, standard database/sql for compatibility
  • Advanced Connection Pooling: Configurable pool sizes, timeouts, and lifecycle management
  • Transaction Management: Support for both GORM and pgx transactions
Migration System
  • Versioned Migrations: golang-migrate integration with embedded SQL files
  • Forward/Backward Support: Up and down migrations for schema evolution
  • Migration Validation: Check migration status and detect dirty states
Vector Search & RAG
  • pgvector Integration: Native PostgreSQL vector similarity search
  • Semantic Search: Find similar entities using cosine similarity
  • Document Chunking: Store and retrieve document chunks for RAG
  • Hybrid Search: Combine full-text and vector search
  • HNSW Indexing: Fast approximate nearest neighbor search
Multi-Layer Caching
  • L1 In-Memory Cache: Ultra-fast local caching with LRU eviction
  • L2 Redis Cache: Distributed caching for scaling across replicas
  • L3 Database Cache: Metadata tracking for cache status
  • Cache Metrics: Hit rates, evictions, and performance tracking
Health Monitoring
  • Comprehensive Health Checks: Connection, query execution, and replication lag
  • Pool Statistics: Track open connections, wait times, and resource usage
  • Table Statistics: Monitor row counts, vacuum status, and index health
  • Extension Checking: Verify required PostgreSQL extensions
Audit & Compliance
  • Comprehensive Audit Logging: Track all data access and modifications
  • PII Access Tracking: Special logging for sensitive data access
  • Compliance Reports: Query audit logs for compliance verification
  • Distributed Tracing Integration: Correlate audit logs with traces
Repository Pattern
  • Type-Safe Models: GORM models with proper relationships
  • Repository Implementations: Catalog, Relationship, Schema repositories
  • Query Builders: Flexible filtering and pagination
  • Preloading Support: Eager load related entities

Installation

go get github.com/click2-run/dictamesh/pkg/database

Quick Start

Basic Setup
import (
    "github.com/click2-run/dictamesh/pkg/database"
    "go.uber.org/zap"
)

// Create configuration
config := database.DefaultConfig()
config.Host = "localhost"
config.Port = 5432
config.User = "dictamesh"
config.Password = "your_password"
config.Database = "metadata_catalog"

// Create logger
logger, _ := zap.NewProduction()

// Create database instance
db, err := database.New(config, logger)
if err != nil {
    log.Fatal(err)
}

// Connect
ctx := context.Background()
if err := db.Connect(ctx); err != nil {
    log.Fatal(err)
}
defer db.Close()
Running Migrations
import "github.com/click2-run/dictamesh/pkg/database/migrations"

// Create migrator
migrator, err := migrations.NewMigrator(db.StdDB(), logger)
if err != nil {
    log.Fatal(err)
}
defer migrator.Close()

// Run all pending migrations
if err := migrator.Up(ctx); err != nil {
    log.Fatal(err)
}

// Check migration status
version, dirty, err := migrator.Version()
fmt.Printf("Current version: %d, Dirty: %v\n", version, dirty)
import (
    "github.com/click2-run/dictamesh/pkg/database"
    "github.com/pgvector/pgvector-go"
)

// Create vector search instance
vs := database.NewVectorSearch(db)

// Store an embedding
embedding := &database.EntityEmbedding{
    CatalogID:          "entity-123",
    EmbeddingModel:     "text-embedding-ada-002",
    EmbeddingVersion:   "v1",
    EmbeddingDimensions: 1536,
    Embedding:          pgvector.NewVector([]float32{0.1, 0.2, ...}),
    SourceText:         "This is the text that was embedded",
}

if err := vs.StoreEmbedding(ctx, embedding); err != nil {
    log.Fatal(err)
}

// Find similar entities
queryVector := pgvector.NewVector([]float32{0.15, 0.25, ...})
similar, err := vs.FindSimilarEntities(
    ctx,
    queryVector,
    "text-embedding-ada-002",
    0.7,  // similarity threshold
    10,   // limit
)

for _, entity := range similar {
    fmt.Printf("Entity: %s, Similarity: %.4f\n", entity.CatalogID, entity.Similarity)
}

// RAG: Find relevant chunks
chunks, err := vs.FindRelevantChunks(
    ctx,
    queryVector,
    "text-embedding-ada-002",
    nil,  // no entity filter
    0.7,  // similarity threshold
    5,    // top 5 chunks
)

for _, chunk := range chunks {
    fmt.Printf("Chunk: %s\nText: %s\n\n", chunk.ChunkID, chunk.ChunkText)
}
Caching
import "github.com/click2-run/dictamesh/pkg/database/cache"

// Create cache
cacheConfig := cache.DefaultConfig()
cacheConfig.RedisURL = "redis://localhost:6379"

cache, err := cache.New(cacheConfig, logger)
if err != nil {
    log.Fatal(err)
}
defer cache.Close()

// Store value
data := []byte("cached data")
if err := cache.Set(ctx, "my-key", data, 5*time.Minute); err != nil {
    log.Error(err)
}

// Retrieve value
value, err := cache.Get(ctx, "my-key")
if err == nil {
    fmt.Printf("Retrieved: %s\n", string(value))
}

// Store JSON
type MyData struct {
    Name string
    Value int
}
if err := cache.SetJSON(ctx, "json-key", MyData{"test", 42}, 10*time.Minute); err != nil {
    log.Error(err)
}

// Get metrics
metrics := cache.GetMetrics()
fmt.Printf("L1 Hits: %d, L2 Hits: %d\n", metrics.L1Hits, metrics.L2Hits)
Health Checks
import "github.com/click2-run/dictamesh/pkg/database/health"

// Create health checker
checker := health.NewChecker(db.Pool(), db.StdDB(), logger)

// Perform health check
result := checker.Check(ctx)
fmt.Printf("Status: %s\nMessage: %s\nResponse Time: %v\n",
    result.Status, result.Message, result.ResponseTime)

// Check specific table (use dictamesh_ prefix)
if err := checker.CheckTable(ctx, "dictamesh_entity_catalog"); err != nil {
    log.Printf("Table check failed: %v", err)
}

// Check extension
if err := checker.CheckExtension(ctx, "vector"); err != nil {
    log.Printf("Extension check failed: %v", err)
}

// Get table statistics (use dictamesh_ prefix)
stats, err := checker.GetTableStats(ctx, "dictamesh_entity_catalog")
if err == nil {
    fmt.Printf("Live tuples: %d\n", stats["live_tuples"])
}
Audit Logging
import "github.com/click2-run/dictamesh/pkg/database/audit"

// Create audit logger
auditConfig := &audit.Config{Enabled: true}
auditLogger := audit.NewLogger(db.Pool(), logger, auditConfig)

// Create audit table
if err := auditLogger.CreateAuditTable(ctx); err != nil {
    log.Fatal(err)
}

// Log an operation
entry := &audit.AuditLog{
    UserID:       "user-123",
    UserEmail:    "user@example.com",
    Operation:    audit.OpUpdate,
    ResourceType: "customer",
    ResourceID:   "cust-456",
    Changes: map[string]interface{}{
        "email": "newemail@example.com",
    },
    Success:    true,
    TraceID:    "trace-789",
}

if err := auditLogger.Log(ctx, entry); err != nil {
    log.Error(err)
}

// Log PII access
if err := auditLogger.LogDataAccess(ctx, "user-123", "customer", "cust-456",
    []string{"ssn", "credit_card"}); err != nil {
    log.Error(err)
}

// Query audit logs
filters := &audit.QueryFilters{
    UserID:     "user-123",
    StartTime:  time.Now().Add(-24 * time.Hour),
    EndTime:    time.Now(),
    Limit:      100,
}

logs, err := auditLogger.Query(ctx, filters)
for _, log := range logs {
    fmt.Printf("%s: %s on %s\n", log.Timestamp, log.Operation, log.ResourceType)
}
Repository Pattern
import (
    "github.com/click2-run/dictamesh/pkg/database/models"
    "github.com/click2-run/dictamesh/pkg/database/repository"
)

// Create repositories
catalogRepo := repository.NewCatalogRepository(db.GORM())
relationshipRepo := repository.NewRelationshipRepository(db.GORM())

// Create entity
entity := &models.EntityCatalog{
    EntityType:     "customer",
    Domain:         "customers",
    SourceSystem:   "directus",
    SourceEntityID: "12345",
    APIBaseURL:     "https://api.example.com",
    APIPathTemplate: "/customers/{id}",
    Status:         "active",
}

if err := catalogRepo.Create(ctx, entity); err != nil {
    log.Fatal(err)
}

// Find entity
found, err := catalogRepo.FindBySource(ctx, "directus", "12345", "customer")
if err != nil {
    log.Fatal(err)
}

// List entities with filters
filters := &repository.CatalogFilters{
    EntityType: "customer",
    Domain:     "customers",
    Status:     "active",
    Limit:      50,
    Offset:     0,
}

entities, err := catalogRepo.List(ctx, filters)
fmt.Printf("Found %d entities\n", len(entities))

Configuration

Database Config
config := &database.Config{
    Host:     "localhost",
    Port:     5432,
    User:     "dictamesh",
    Password: "password",
    Database: "metadata_catalog",
    SSLMode:  "prefer",

    // Connection pool
    MaxOpenConns:    25,
    MaxIdleConns:    10,
    ConnMaxLifetime: 30 * time.Minute,
    ConnMaxIdleTime: 10 * time.Minute,

    // Performance
    StatementTimeout: 30 * time.Second,
    IdleInTxTimeout:  60 * time.Second,

    // Features
    EnableMigrations:   true,
    EnableVectorSearch: true,
    EnableAuditLog:     true,

    // Observability
    EnableMetrics: true,
    EnableTracing: true,
    LogLevel:      "info",
}

Schema

The database package includes comprehensive schema migrations:

  • 000001_initial_schema.up.sql: Core metadata catalog tables
  • 000002_add_vector_search.up.sql: Vector embeddings and RAG support
Tables

IMPORTANT: All tables use the dictamesh_ prefix for namespace isolation.

  • dictamesh_entity_catalog: Registry of all entities
  • dictamesh_entity_relationships: Cross-system relationships with temporal validity
  • dictamesh_schemas: Versioned entity schemas
  • dictamesh_event_log: Immutable audit trail
  • dictamesh_data_lineage: Data flow tracking
  • dictamesh_cache_status: Cache freshness tracking
  • dictamesh_entity_embeddings: Vector embeddings for semantic search
  • dictamesh_document_chunks: Document chunks for RAG
  • dictamesh_audit_logs: Comprehensive audit logging

Performance Tips

  1. Use pgx for high-performance queries: Direct pool access for bulk operations
  2. Enable caching: Reduce database load with multi-layer caching
  3. Batch operations: Use transactions for multiple operations
  4. Optimize indexes: Vector search uses HNSW for fast retrieval
  5. Monitor connections: Track pool stats and adjust configuration

Security

  • Connection encryption: SSL/TLS support
  • Prepared statements: Protection against SQL injection
  • Audit logging: Track all sensitive data access
  • PII tracking: Special handling for personal information

License

AGPL-3.0-or-later - Copyright (C) 2025 Controle Digital Ltda

Documentation

Overview

Package database provides comprehensive database infrastructure for DictaMesh. It includes connection pooling, migrations, ORM, vector search, caching, and audit logging.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Config

type Config struct {
	// Connection settings
	Host     string
	Port     int
	User     string
	Password string
	Database string
	SSLMode  string

	// Connection pool settings
	MaxOpenConns    int
	MaxIdleConns    int
	ConnMaxLifetime time.Duration
	ConnMaxIdleTime time.Duration

	// Performance settings
	StatementTimeout time.Duration
	IdleInTxTimeout  time.Duration

	// Feature flags
	EnableMigrations   bool
	EnableVectorSearch bool
	EnableAuditLog     bool

	// Observability
	EnableMetrics bool
	EnableTracing bool
	LogLevel      string
}

Config represents database configuration

func DefaultConfig

func DefaultConfig() *Config

DefaultConfig returns a production-ready default configuration

func (*Config) DSN

func (c *Config) DSN() string

DSN returns the PostgreSQL connection string

func (*Config) Validate

func (c *Config) Validate() error

Validate checks if the configuration is valid

type Database

type Database struct {
	// contains filtered or unexported fields
}

Database represents the main database connection manager

func New

func New(config *Config, logger *zap.Logger) (*Database, error)

New creates a new database instance

func (*Database) Close

func (db *Database) Close() error

Close closes all database connections

func (*Database) Connect

func (db *Database) Connect(ctx context.Context) error

Connect establishes database connections

func (*Database) GORM

func (db *Database) GORM() *gorm.DB

GORM returns the GORM instance for ORM operations

func (*Database) GetMetrics

func (db *Database) GetMetrics() *Metrics

GetMetrics returns current performance metrics

func (*Database) Ping

func (db *Database) Ping(ctx context.Context) error

Ping checks if the database is reachable

func (*Database) Pool

func (db *Database) Pool() *pgxpool.Pool

Pool returns the pgx connection pool for high-performance queries

func (*Database) Stats

func (db *Database) Stats() sql.DBStats

Stats returns database statistics

func (*Database) StdDB

func (db *Database) StdDB() *sql.DB

StdDB returns the standard database/sql instance

func (*Database) WithPgxTransaction

func (db *Database) WithPgxTransaction(ctx context.Context, fn func(pgx.Tx) error) error

WithPgxTransaction executes a function within a pgx transaction

func (*Database) WithTransaction

func (db *Database) WithTransaction(ctx context.Context, fn func(*gorm.DB) error) error

WithTransaction executes a function within a database transaction

type DocumentChunk

type DocumentChunk struct {
	ID               string
	CatalogID        string
	ChunkIndex       int
	ChunkText        string
	ChunkTokens      int
	EmbeddingModel   string
	Embedding        pgvector.Vector
	PrecedingContext string
	FollowingContext string
	Metadata         map[string]interface{}
}

DocumentChunk represents a chunked document for RAG

type EmbeddingModel

type EmbeddingModel struct {
	Name       string
	Version    string
	Dimensions int
}

EmbeddingModel represents an embedding model configuration

type EntityEmbedding

type EntityEmbedding struct {
	ID                  string
	CatalogID           string
	EmbeddingModel      string
	EmbeddingVersion    string
	EmbeddingDimensions int
	Embedding           pgvector.Vector
	SourceText          string
	SourceFields        map[string]interface{}
	Metadata            map[string]interface{}
}

EntityEmbedding represents a vector embedding of an entity

type HybridSearchResult

type HybridSearchResult struct {
	CatalogID        string
	CombinedScore    float64
	TextRank         float64
	VectorSimilarity float64
	SourceText       string
}

HybridSearchResult represents a result from hybrid search

type Metrics

type Metrics struct {
	QueryCount       int64
	QueryErrors      int64
	CacheHits        int64
	CacheMisses      int64
	ConnectionsOpen  int32
	ConnectionsIdle  int32
	AvgQueryDuration time.Duration
}

Metrics tracks database performance metrics

type RelevantChunk

type RelevantChunk struct {
	ChunkID          string
	CatalogID        string
	ChunkText        string
	ChunkIndex       int
	PrecedingContext string
	FollowingContext string
	Similarity       float64
	Metadata         map[string]interface{}
}

RelevantChunk represents a relevant document chunk for RAG

type SimilarEntity

type SimilarEntity struct {
	CatalogID  string
	Similarity float64
	SourceText string
	Metadata   map[string]interface{}
}

SimilarEntity represents a search result with similarity score

type VectorSearch

type VectorSearch struct {
	// contains filtered or unexported fields
}

VectorSearch provides vector similarity search capabilities

func NewVectorSearch

func NewVectorSearch(db *Database) *VectorSearch

NewVectorSearch creates a new vector search instance

func (*VectorSearch) BatchStoreChunks

func (vs *VectorSearch) BatchStoreChunks(ctx context.Context, chunks []DocumentChunk) error

BatchStoreChunks stores multiple document chunks in a transaction

func (*VectorSearch) DeleteDocumentChunks

func (vs *VectorSearch) DeleteDocumentChunks(ctx context.Context, catalogID string) error

DeleteDocumentChunks deletes all chunks for a catalog entry

func (*VectorSearch) DeleteEmbeddings

func (vs *VectorSearch) DeleteEmbeddings(ctx context.Context, catalogID string) error

DeleteEmbeddings deletes all embeddings for a catalog entry

func (*VectorSearch) FindRelevantChunks

func (vs *VectorSearch) FindRelevantChunks(
	ctx context.Context,
	queryEmbedding pgvector.Vector,
	modelName string,
	catalogID *string,
	similarityThreshold float64,
	limit int,
) ([]RelevantChunk, error)

FindRelevantChunks finds relevant document chunks for RAG

func (*VectorSearch) FindSimilarEntities

func (vs *VectorSearch) FindSimilarEntities(
	ctx context.Context,
	queryEmbedding pgvector.Vector,
	modelName string,
	similarityThreshold float64,
	limit int,
) ([]SimilarEntity, error)

FindSimilarEntities finds entities similar to the query embedding

func (*VectorSearch) HybridSearch

func (vs *VectorSearch) HybridSearch(
	ctx context.Context,
	queryText string,
	queryEmbedding pgvector.Vector,
	modelName string,
	textWeight float64,
	vectorWeight float64,
	limit int,
) ([]HybridSearchResult, error)

HybridSearch performs combined full-text and vector search

func (*VectorSearch) StoreDocumentChunk

func (vs *VectorSearch) StoreDocumentChunk(ctx context.Context, chunk *DocumentChunk) error

StoreDocumentChunk stores a document chunk with embedding

func (*VectorSearch) StoreEmbedding

func (vs *VectorSearch) StoreEmbedding(ctx context.Context, embedding *EntityEmbedding) error

StoreEmbedding stores an entity embedding

Directories

Path Synopsis
Package audit provides comprehensive audit logging and compliance tracking
Package audit provides comprehensive audit logging and compliance tracking
Package cache provides multi-layer caching with Redis integration
Package cache provides multi-layer caching with Redis integration
Package health provides database health checking and monitoring
Package health provides database health checking and monitoring
Package models provides database models for the metadata catalog
Package models provides database models for the metadata catalog
Package repository provides repository pattern implementations for database access
Package repository provides repository pattern implementations for database access

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL