package main import ( "crypto/rand" "database/sql" "encoding/json" "log" "net/http" "os" "strconv" "strings" "time" _ "modernc.org/sqlite" ) var ( db *sql.DB hostname string dbFile string ) func main() { hostname, _ = os.Hostname() dbPath := os.Getenv("DB_PATH") if dbPath == "" { dbPath = "/mnt/data/poc.db" } dbFile = dbPath var err error db, err = sql.Open("sqlite", "file:"+dbPath+"?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)&_pragma=synchronous(NORMAL)") if err != nil { log.Fatal(err) } if _, err := db.Exec(`CREATE TABLE IF NOT EXISTS entries ( id INTEGER PRIMARY KEY AUTOINCREMENT, container TEXT NOT NULL, msg TEXT NOT NULL, created_at TEXT NOT NULL, payload BLOB )`); err != nil { log.Fatal(err) } // the payload column was added after the first deploy, older DB files need it retrofitted if _, err := db.Exec(`ALTER TABLE entries ADD COLUMN payload BLOB`); err != nil && !strings.Contains(err.Error(), "duplicate column") { log.Fatal(err) } http.HandleFunc("/", handleStatus) http.HandleFunc("/write", handleWrite) http.HandleFunc("/read", handleRead) http.HandleFunc("/checkpoint", handleCheckpoint) log.Printf("listening on :8080, db=%s, container=%s", dbPath, hostname) log.Fatal(http.ListenAndServe(":8080", nil)) } func respond(w http.ResponseWriter, status int, data any) { w.Header().Set("Content-Type", "application/json") w.WriteHeader(status) json.NewEncoder(w).Encode(data) } func handleStatus(w http.ResponseWriter, r *http.Request) { var count int64 var journalMode string if err := db.QueryRow("SELECT COUNT(*) FROM entries").Scan(&count); err != nil { respond(w, 500, map[string]any{"container": hostname, "error": err.Error()}) return } db.QueryRow("PRAGMA journal_mode").Scan(&journalMode) fileSize := func(path string) int64 { info, err := os.Stat(path) if err != nil { return 0 } return info.Size() } respond(w, 200, map[string]any{ "container": hostname, "entries": count, "journalMode": journalMode, "dbSizeBytes": fileSize(dbFile), "walSizeBytes": fileSize(dbFile + "-wal"), }) } func handleWrite(w http.ResponseWriter, r *http.Request) { msg := r.URL.Query().Get("msg") if msg == "" { msg = "hello" } intParam := func(key string, def, max int) int { v, err := strconv.Atoi(r.URL.Query().Get(key)) if err != nil || v < 1 { return def } return min(v, max) } n := intParam("n", 1, 10000) size := intParam("size", 0, 1<<20) tx, err := db.Begin() if err != nil { respond(w, 500, map[string]any{"container": hostname, "error": err.Error()}) return } defer tx.Rollback() stmt, err := tx.Prepare("INSERT INTO entries (container, msg, created_at, payload) VALUES (?, ?, ?, ?)") if err != nil { respond(w, 500, map[string]any{"container": hostname, "error": err.Error()}) return } defer stmt.Close() var lastId int64 now := time.Now().UTC().Format(time.RFC3339Nano) for i := 0; i < n; i++ { var payload []byte if size > 0 { payload = make([]byte, size) rand.Read(payload) } res, err := stmt.Exec(hostname, msg, now, payload) if err != nil { respond(w, 500, map[string]any{"container": hostname, "error": err.Error()}) return } lastId, _ = res.LastInsertId() } if err := tx.Commit(); err != nil { respond(w, 500, map[string]any{"container": hostname, "error": err.Error()}) return } respond(w, 200, map[string]any{"container": hostname, "inserted": n, "lastId": lastId}) } func handleCheckpoint(w http.ResponseWriter, r *http.Request) { var busy, logFrames, checkpointed int64 if err := db.QueryRow("PRAGMA wal_checkpoint(TRUNCATE)").Scan(&busy, &logFrames, &checkpointed); err != nil { respond(w, 500, map[string]any{"container": hostname, "error": err.Error()}) return } respond(w, 200, map[string]any{"container": hostname, "busy": busy, "logFrames": logFrames, "checkpointed": checkpointed}) } func handleRead(w http.ResponseWriter, r *http.Request) { var count int64 if err := db.QueryRow("SELECT COUNT(*) FROM entries").Scan(&count); err != nil { respond(w, 500, map[string]any{"container": hostname, "error": err.Error()}) return } rows, err := db.Query("SELECT id, container, msg, created_at FROM entries ORDER BY id DESC LIMIT 5") if err != nil { respond(w, 500, map[string]any{"container": hostname, "error": err.Error()}) return } defer rows.Close() type entry struct { Id int64 `json:"id"` Container string `json:"container"` Msg string `json:"msg"` CreatedAt string `json:"createdAt"` } last := []entry{} for rows.Next() { var e entry if err := rows.Scan(&e.Id, &e.Container, &e.Msg, &e.CreatedAt); err != nil { respond(w, 500, map[string]any{"container": hostname, "error": err.Error()}) return } last = append(last, e) } respond(w, 200, map[string]any{"container": hostname, "entries": count, "last": last}) }