179 lines
4.8 KiB
Go
179 lines
4.8 KiB
Go
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})
|
|
}
|