Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| cbc5022e0d | |||
| 74b32800f9 | |||
| cb13839157 | |||
| f833498b65 | |||
| 6064ddd856 | |||
| dc56261b58 |
@@ -8,7 +8,7 @@ Create a PR on Gitea, wait for CI, and squash-merge it. Push code to both remote
|
|||||||
5. Create a Gitea PR: `tea pr create --head <branch> --base main`. Reference related issues with `#X`. Only use `Closes #X` if the PR fully resolves the issue.
|
5. Create a Gitea PR: `tea pr create --head <branch> --base main`. Reference related issues with `#X`. Only use `Closes #X` if the PR fully resolves the issue.
|
||||||
6. Wait for CI to pass: poll Gitea CI status. If CI fails, report the failure and stop — do not merge.
|
6. Wait for CI to pass: poll Gitea CI status. If CI fails, report the failure and stop — do not merge.
|
||||||
7. Once CI passes, squash-merge on Gitea: `tea pr merge <index> --style squash` with a clean, semantic commit message including the PR number. No Claude attribution lines.
|
7. Once CI passes, squash-merge on Gitea: `tea pr merge <index> --style squash` with a clean, semantic commit message including the PR number. No Claude attribution lines.
|
||||||
8. Update local main and push to both remotes: `git checkout main && git pull gitea main && git push origin main`.
|
8. Update local main and push to both remotes. If in a worktree, `main` is checked out in the primary tree, so run from there: `cd <primary-worktree> && git pull gitea main && git push origin main` (the primary worktree path is the repo root without `.claude/worktrees/…`). If not in a worktree: `git checkout main && git pull gitea main && git push origin main`.
|
||||||
9. Clean up remote branches: `git push gitea --delete <branch> && git push origin --delete <branch>`.
|
9. Clean up remote branches: `git push gitea --delete <branch> && git push origin --delete <branch>`.
|
||||||
10. Prune refs: `git remote prune gitea && git remote prune origin`.
|
10. Prune refs: `git remote prune gitea && git remote prune origin`.
|
||||||
11. Report the merged PR URL.
|
11. Report the merged PR URL.
|
||||||
|
|||||||
18
.gitignore
vendored
18
.gitignore
vendored
@@ -37,21 +37,15 @@ go.work.sum
|
|||||||
# Air artifacts
|
# Air artifacts
|
||||||
*tmp/
|
*tmp/
|
||||||
|
|
||||||
# binaries
|
# Example binaries and data files
|
||||||
internal/examples/chatroom/chatroom
|
internal/examples/*/[a-z]*[!.go]
|
||||||
internal/examples/counter/counter
|
internal/examples/shakespeare/shake.db
|
||||||
internal/examples/countercomp/countercomp
|
|
||||||
internal/examples/greeter/greeter
|
|
||||||
internal/examples/livereload/livereload
|
|
||||||
internal/examples/picocss/picocss
|
|
||||||
internal/examples/plugins/plugins
|
|
||||||
internal/examples/realtimechart/realtimechart
|
|
||||||
internal/examples/shakespeare/shakespeare
|
|
||||||
internal/examples/nats-chatroom/nats-chatroom
|
|
||||||
/nats-chatroom
|
|
||||||
|
|
||||||
# NATS data directory
|
# NATS data directory
|
||||||
data/
|
data/
|
||||||
|
|
||||||
|
# Standalone experiments
|
||||||
|
nats-chatroom
|
||||||
|
|
||||||
# Claude Code worktrees
|
# Claude Code worktrees
|
||||||
.claude/worktrees/
|
.claude/worktrees/
|
||||||
|
|||||||
@@ -74,14 +74,13 @@ func main() {
|
|||||||
- **Plugin system** — `func(v *V)` hooks for integrating CSS/JS libraries
|
- **Plugin system** — `func(v *V)` hooks for integrating CSS/JS libraries
|
||||||
- **Structured logging** — zerolog with configurable levels; console output in dev, JSON in production
|
- **Structured logging** — zerolog with configurable levels; console output in dev, JSON in production
|
||||||
- **Graceful shutdown** — listens for SIGINT/SIGTERM, drains contexts, closes pub/sub
|
- **Graceful shutdown** — listens for SIGINT/SIGTERM, drains contexts, closes pub/sub
|
||||||
- **Context lifecycle** — background reaper cleans up disconnected contexts; configurable TTL
|
|
||||||
- **HTML DSL** — the `h` package provides type-safe Go-native HTML composition
|
- **HTML DSL** — the `h` package provides type-safe Go-native HTML composition
|
||||||
|
|
||||||
## Examples
|
## Examples
|
||||||
|
|
||||||
The `internal/examples/` directory contains 14 runnable examples:
|
The `internal/examples/` directory contains 19 runnable examples:
|
||||||
|
|
||||||
`chatroom` · `counter` · `countercomp` · `greeter` · `keyboard` · `livereload` · `nats-chatroom` · `pathparams` · `picocss` · `plugins` · `pubsub-crud` · `realtimechart` · `session` · `shakespeare`
|
`chatroom` · `counter` · `countercomp` · `effectspike` · `greeter` · `keyboard` · `livereload` · `maplibre` · `middleware` · `nats-chatroom` · `pathparams` · `picocss` · `plugins` · `pubsub-crud` · `realtimechart` · `session` · `shakespeare` · `signup` · `spa`
|
||||||
|
|
||||||
## Experimental
|
## Experimental
|
||||||
|
|
||||||
|
|||||||
10
ci-check.sh
10
ci-check.sh
@@ -6,9 +6,13 @@ set -o pipefail
|
|||||||
ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
|
ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
|
||||||
cd "$ROOT"
|
cd "$ROOT"
|
||||||
|
|
||||||
echo "== CI: Format code =="
|
echo "== CI: Check formatting =="
|
||||||
go fmt ./...
|
if [ -n "$(gofmt -l .)" ]; then
|
||||||
echo "OK: formatting complete"
|
echo "ERROR: files not formatted:"
|
||||||
|
gofmt -l .
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
echo "OK: all files formatted"
|
||||||
|
|
||||||
echo "== CI: Run go vet =="
|
echo "== CI: Run go vet =="
|
||||||
if ! go vet ./...; then
|
if ! go vet ./...; then
|
||||||
|
|||||||
@@ -1,7 +1,6 @@
|
|||||||
package via
|
package via
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
@@ -51,5 +50,5 @@ func (s *computedSignal) recompute() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (s *computedSignal) patchValue() string {
|
func (s *computedSignal) patchValue() string {
|
||||||
return fmt.Sprintf("%v", s.lastVal)
|
return s.lastVal
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,8 +1,6 @@
|
|||||||
package via
|
package via
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/alexedwards/scs/v2"
|
"github.com/alexedwards/scs/v2"
|
||||||
"github.com/rs/zerolog"
|
"github.com/rs/zerolog"
|
||||||
)
|
)
|
||||||
@@ -62,16 +60,6 @@ type Options struct {
|
|||||||
// the embedded NATS server. Ignored when a custom PubSub is configured.
|
// the embedded NATS server. Ignored when a custom PubSub is configured.
|
||||||
Streams []StreamConfig
|
Streams []StreamConfig
|
||||||
|
|
||||||
// ContextSuspendAfter is the time a context may be disconnected before
|
|
||||||
// the reaper suspends it (frees page resources but keeps the context
|
|
||||||
// shell for seamless re-init on reconnect). Default: 15m.
|
|
||||||
ContextSuspendAfter time.Duration
|
|
||||||
|
|
||||||
// ContextTTL is the maximum time a context may exist without an SSE
|
|
||||||
// connection before the background reaper fully disposes it.
|
|
||||||
// Default: 1h. Negative value disables the reaper.
|
|
||||||
ContextTTL time.Duration
|
|
||||||
|
|
||||||
// ActionRateLimit configures the default token-bucket rate limiter for
|
// ActionRateLimit configures the default token-bucket rate limiter for
|
||||||
// action endpoints. Zero values use built-in defaults (10 req/s, burst 20).
|
// action endpoints. Zero values use built-in defaults (10 req/s, burst 20).
|
||||||
// Set Rate to -1 to disable rate limiting entirely.
|
// Set Rate to -1 to disable rate limiting entirely.
|
||||||
|
|||||||
@@ -42,8 +42,6 @@ type Context struct {
|
|||||||
createdAt time.Time
|
createdAt time.Time
|
||||||
sseConnected atomic.Bool
|
sseConnected atomic.Bool
|
||||||
sseDisconnectedAt atomic.Pointer[time.Time]
|
sseDisconnectedAt atomic.Pointer[time.Time]
|
||||||
lastSeenAt atomic.Pointer[time.Time]
|
|
||||||
suspended atomic.Bool
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// View defines the UI rendered by this context.
|
// View defines the UI rendered by this context.
|
||||||
@@ -444,13 +442,6 @@ func (c *Context) resetPageState() {
|
|||||||
c.mu.Unlock()
|
c.mu.Unlock()
|
||||||
}
|
}
|
||||||
|
|
||||||
// suspend frees page-scoped resources while keeping the context shell alive
|
|
||||||
// in the registry for seamless re-init on reconnect.
|
|
||||||
func (c *Context) suspend() {
|
|
||||||
c.resetPageState()
|
|
||||||
c.suspended.Store(true)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Navigate performs an SPA navigation to the given path. It resets page state,
|
// Navigate performs an SPA navigation to the given path. It resets page state,
|
||||||
// runs the target page's init function (with middleware), and pushes the new
|
// runs the target page's init function (with middleware), and pushes the new
|
||||||
// view over the existing SSE connection with a view transition animation.
|
// view over the existing SSE connection with a view transition animation.
|
||||||
|
|||||||
@@ -66,7 +66,6 @@ v.Config(via.Options{
|
|||||||
Plugins: []via.Plugin{MyPlugin},
|
Plugins: []via.Plugin{MyPlugin},
|
||||||
SessionManager: sm,
|
SessionManager: sm,
|
||||||
PubSub: customBackend,
|
PubSub: customBackend,
|
||||||
ContextTTL: 60 * time.Second,
|
|
||||||
ActionRateLimit: via.RateLimitConfig{Rate: 20, Burst: 40},
|
ActionRateLimit: via.RateLimitConfig{Rate: 20, Burst: 40},
|
||||||
})
|
})
|
||||||
```
|
```
|
||||||
@@ -83,7 +82,6 @@ v.Config(via.Options{
|
|||||||
| `DatastarContent` | (embedded) | Custom Datastar JS bytes |
|
| `DatastarContent` | (embedded) | Custom Datastar JS bytes |
|
||||||
| `DatastarPath` | `"/_datastar.js"` | URL path for the Datastar script |
|
| `DatastarPath` | `"/_datastar.js"` | URL path for the Datastar script |
|
||||||
| `PubSub` | embedded NATS | Custom PubSub backend. Replaces the default NATS. See [PubSub and Sessions](pubsub-and-sessions.md) |
|
| `PubSub` | embedded NATS | Custom PubSub backend. Replaces the default NATS. See [PubSub and Sessions](pubsub-and-sessions.md) |
|
||||||
| `ContextTTL` | `30s` | Max time a context survives without an SSE connection before cleanup. Negative value disables the reaper |
|
|
||||||
| `ActionRateLimit` | `10 req/s, burst 20` | Default token-bucket rate limiter for action endpoints. Rate of `-1` disables limiting |
|
| `ActionRateLimit` | `10 req/s, burst 20` | Default token-bucket rate limiter for action endpoints. Rate of `-1` disables limiting |
|
||||||
|
|
||||||
## Static Files
|
## Static Files
|
||||||
|
|||||||
@@ -14,7 +14,7 @@ Browser hits page → new Context created → init function runs → HTML render
|
|||||||
action fires → signals injected from browser → handler runs → Sync() → DOM patched
|
action fires → signals injected from browser → handler runs → Sync() → DOM patched
|
||||||
```
|
```
|
||||||
|
|
||||||
The context is disposed when the SSE connection closes (tab close, navigation away, network loss). A background reaper also cleans up contexts that never establish an SSE connection within `ContextTTL` (default 30s).
|
The context lives until the browser tab closes (detected via a `beforeunload` beacon) or the server shuts down. There is no background reaper — contexts persist across temporary SSE disconnections so backgrounded tabs resume seamlessly.
|
||||||
|
|
||||||
During [SPA navigation](routing-and-navigation.md#spa-navigation), the context itself survives — only page-level state (signals, actions, fields, intervals, subscriptions) is reset. The SSE connection persists.
|
During [SPA navigation](routing-and-navigation.md#spa-navigation), the context itself survives — only page-level state (signals, actions, fields, intervals, subscriptions) is reset. The SSE connection persists.
|
||||||
|
|
||||||
@@ -94,7 +94,7 @@ Available triggers:
|
|||||||
|--------|-------|-------|
|
|--------|-------|-------|
|
||||||
| `OnClick()` | `click` | |
|
| `OnClick()` | `click` | |
|
||||||
| `OnDblClick()` | `dblclick` | |
|
| `OnDblClick()` | `dblclick` | |
|
||||||
| `OnChange()` | `change` | 200ms debounce |
|
| `OnChange()` | `change` | |
|
||||||
| `OnInput()` | `input` | No debounce |
|
| `OnInput()` | `input` | No debounce |
|
||||||
| `OnSubmit()` | `submit` | |
|
| `OnSubmit()` | `submit` | |
|
||||||
| `OnKeyDown(key)` | `keydown` | Filtered by key name (e.g. `"Enter"`, `"Escape"`) |
|
| `OnKeyDown(key)` | `keydown` | Filtered by key name (e.g. `"Enter"`, `"Escape"`) |
|
||||||
|
|||||||
2
h/h.go
2
h/h.go
@@ -5,7 +5,7 @@
|
|||||||
//
|
//
|
||||||
// h.Div(
|
// h.Div(
|
||||||
// h.H1(h.Text("Hello, Via")),
|
// h.H1(h.Text("Hello, Via")),
|
||||||
// h.P(h.Text("Pure Go. No tmplates.")),
|
// h.P(h.Text("Pure Go. No templates.")),
|
||||||
// )
|
// )
|
||||||
package h
|
package h
|
||||||
|
|
||||||
|
|||||||
@@ -2,7 +2,6 @@ package main
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"github.com/ryanhamamura/via"
|
"github.com/ryanhamamura/via"
|
||||||
// "github.com/go-via/via-plugin-picocss/picocss"
|
|
||||||
"github.com/ryanhamamura/via/h"
|
"github.com/ryanhamamura/via/h"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -15,9 +14,6 @@ func main() {
|
|||||||
DocumentTitle: "Live Reload Demo",
|
DocumentTitle: "Live Reload Demo",
|
||||||
DevMode: true,
|
DevMode: true,
|
||||||
LogLevel: via.LogLevelDebug,
|
LogLevel: via.LogLevelDebug,
|
||||||
Plugins: []via.Plugin{
|
|
||||||
// picocss.Default
|
|
||||||
},
|
|
||||||
})
|
})
|
||||||
|
|
||||||
v.Page("/", func(c *via.Context) {
|
v.Page("/", func(c *via.Context) {
|
||||||
|
|||||||
@@ -2,6 +2,8 @@ package main
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"math/rand"
|
"math/rand"
|
||||||
|
"strconv"
|
||||||
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ryanhamamura/via"
|
"github.com/ryanhamamura/via"
|
||||||
@@ -9,6 +11,19 @@ import (
|
|||||||
"github.com/ryanhamamura/via/maplibre"
|
"github.com/ryanhamamura/via/maplibre"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
type posMsg struct {
|
||||||
|
Lng float64 `json:"lng"`
|
||||||
|
Lat float64 `json:"lat"`
|
||||||
|
}
|
||||||
|
|
||||||
|
var (
|
||||||
|
vehicleOnce sync.Once
|
||||||
|
vehicle struct {
|
||||||
|
mu sync.RWMutex
|
||||||
|
lng, lat float64
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
v := via.New()
|
v := via.New()
|
||||||
v.Config(via.Options{
|
v.Config(via.Options{
|
||||||
@@ -18,6 +33,21 @@ func main() {
|
|||||||
Plugins: []via.Plugin{maplibre.Plugin},
|
Plugins: []via.Plugin{maplibre.Plugin},
|
||||||
})
|
})
|
||||||
|
|
||||||
|
// Single goroutine moves the vehicle — all clients read the same position.
|
||||||
|
vehicle.lng = -122.43
|
||||||
|
vehicle.lat = 37.77
|
||||||
|
vehicleOnce.Do(func() {
|
||||||
|
go func() {
|
||||||
|
for {
|
||||||
|
time.Sleep(time.Second)
|
||||||
|
vehicle.mu.Lock()
|
||||||
|
vehicle.lng = -122.43 + (rand.Float64()-0.5)*0.02
|
||||||
|
vehicle.lat = 37.77 + (rand.Float64()-0.5)*0.02
|
||||||
|
vehicle.mu.Unlock()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
})
|
||||||
|
|
||||||
v.Page("/", func(c *via.Context) {
|
v.Page("/", func(c *via.Context) {
|
||||||
m := maplibre.New(c, maplibre.Options{
|
m := maplibre.New(c, maplibre.Options{
|
||||||
Style: "https://demotiles.maplibre.org/style.json",
|
Style: "https://demotiles.maplibre.org/style.json",
|
||||||
@@ -45,7 +75,7 @@ func main() {
|
|||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
|
||||||
// Signal-backed marker — server pushes position updates
|
// Purple vehicle marker — reads shared Go state
|
||||||
vehicleLng := c.Signal(-122.43)
|
vehicleLng := c.Signal(-122.43)
|
||||||
vehicleLat := c.Signal(37.77)
|
vehicleLat := c.Signal(37.77)
|
||||||
|
|
||||||
@@ -56,12 +86,37 @@ func main() {
|
|||||||
})
|
})
|
||||||
|
|
||||||
c.OnInterval(time.Second, func() {
|
c.OnInterval(time.Second, func() {
|
||||||
vehicleLng.SetValue(-122.43 + (rand.Float64()-0.5)*0.02)
|
vehicle.mu.RLock()
|
||||||
vehicleLat.SetValue(37.77 + (rand.Float64()-0.5)*0.02)
|
lng, lat := vehicle.lng, vehicle.lat
|
||||||
|
vehicle.mu.RUnlock()
|
||||||
|
vehicleLng.SetValue(lng)
|
||||||
|
vehicleLat.SetValue(lat)
|
||||||
c.SyncSignals()
|
c.SyncSignals()
|
||||||
})
|
})
|
||||||
|
|
||||||
// Draggable marker — user drags, signals update
|
// Yellow click marker — synced across clients via PubSub
|
||||||
|
clickLng := c.Signal(-122.4194)
|
||||||
|
clickLat := c.Signal(37.7749)
|
||||||
|
|
||||||
|
m.AddMarker("clicked", maplibre.Marker{
|
||||||
|
LngSignal: clickLng,
|
||||||
|
LatSignal: clickLat,
|
||||||
|
Color: "#f39c12",
|
||||||
|
})
|
||||||
|
|
||||||
|
via.Subscribe(c, "map.click", func(msg posMsg) {
|
||||||
|
clickLng.SetValue(msg.Lng)
|
||||||
|
clickLat.SetValue(msg.Lat)
|
||||||
|
c.SyncSignals()
|
||||||
|
})
|
||||||
|
|
||||||
|
click := m.OnClick()
|
||||||
|
handleClick := c.Action(func() {
|
||||||
|
e := click.Data()
|
||||||
|
via.Publish(c, "map.click", posMsg{Lng: e.LngLat.Lng, Lat: e.LngLat.Lat})
|
||||||
|
})
|
||||||
|
|
||||||
|
// Blue draggable pin — synced across clients via PubSub
|
||||||
pinLng := c.Signal(-122.41)
|
pinLng := c.Signal(-122.41)
|
||||||
pinLat := c.Signal(37.78)
|
pinLat := c.Signal(37.78)
|
||||||
|
|
||||||
@@ -72,14 +127,16 @@ func main() {
|
|||||||
Draggable: true,
|
Draggable: true,
|
||||||
})
|
})
|
||||||
|
|
||||||
// Click event — click to place a marker
|
via.Subscribe(c, "map.pin", func(msg posMsg) {
|
||||||
click := m.OnClick()
|
pinLng.SetValue(msg.Lng)
|
||||||
handleClick := c.Action(func() {
|
pinLat.SetValue(msg.Lat)
|
||||||
e := click.Data()
|
c.SyncSignals()
|
||||||
m.AddMarker("clicked", maplibre.Marker{
|
|
||||||
LngLat: e.LngLat,
|
|
||||||
Color: "#f39c12",
|
|
||||||
})
|
})
|
||||||
|
|
||||||
|
handlePinDrag := c.Action(func() {
|
||||||
|
lng, _ := strconv.ParseFloat(pinLng.String(), 64)
|
||||||
|
lat, _ := strconv.ParseFloat(pinLat.String(), 64)
|
||||||
|
via.Publish(c, "map.pin", posMsg{Lng: lng, Lat: lat})
|
||||||
})
|
})
|
||||||
|
|
||||||
// GeoJSON polygon source + fill layer
|
// GeoJSON polygon source + fill layer
|
||||||
@@ -111,7 +168,7 @@ func main() {
|
|||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
|
||||||
// FlyTo actions using CameraOptions
|
// FlyTo actions
|
||||||
zoom14 := 14.0
|
zoom14 := 14.0
|
||||||
flyToSF := c.Action(func() {
|
flyToSF := c.Action(func() {
|
||||||
m.FlyTo(maplibre.CameraOptions{
|
m.FlyTo(maplibre.CameraOptions{
|
||||||
@@ -134,6 +191,7 @@ func main() {
|
|||||||
h.H1(h.Text("MapLibre GL Example")),
|
h.H1(h.Text("MapLibre GL Example")),
|
||||||
m.Element(
|
m.Element(
|
||||||
click.Input(handleClick.OnInput()),
|
click.Input(handleClick.OnInput()),
|
||||||
|
h.Input(h.Type("hidden"), pinLng.Bind(), handlePinDrag.OnInput()),
|
||||||
),
|
),
|
||||||
h.Div(h.Attr("style", "margin-top:1rem;display:flex;gap:0.5rem;flex-wrap:wrap"),
|
h.Div(h.Attr("style", "margin-top:1rem;display:flex;gap:0.5rem;flex-wrap:wrap"),
|
||||||
h.Button(h.Text("Fly to San Francisco"), flyToSF.OnClick()),
|
h.Button(h.Text("Fly to San Francisco"), flyToSF.OnClick()),
|
||||||
@@ -142,6 +200,7 @@ func main() {
|
|||||||
h.Div(h.Attr("style", "margin-top:0.5rem;font-size:0.9rem"),
|
h.Div(h.Attr("style", "margin-top:0.5rem;font-size:0.9rem"),
|
||||||
h.P(h.Text("Zoom: "), m.Zoom.Text()),
|
h.P(h.Text("Zoom: "), m.Zoom.Text()),
|
||||||
h.P(h.Text("Center: "), m.CenterLng.Text(), h.Text(", "), m.CenterLat.Text()),
|
h.P(h.Text("Center: "), m.CenterLng.Text(), h.Text(", "), m.CenterLat.Text()),
|
||||||
|
h.P(h.Text("Click: "), clickLng.Text(), h.Text(", "), clickLat.Text()),
|
||||||
h.P(h.Text("Vehicle: "), vehicleLng.Text(), h.Text(", "), vehicleLat.Text()),
|
h.P(h.Text("Vehicle: "), vehicleLng.Text(), h.Text(", "), vehicleLat.Text()),
|
||||||
h.P(h.Text("Draggable Pin: "), pinLng.Text(), h.Text(", "), pinLat.Text()),
|
h.P(h.Text("Draggable Pin: "), pinLng.Text(), h.Text(", "), pinLat.Text()),
|
||||||
),
|
),
|
||||||
|
|||||||
@@ -4,17 +4,12 @@ import (
|
|||||||
"strconv"
|
"strconv"
|
||||||
|
|
||||||
"github.com/ryanhamamura/via"
|
"github.com/ryanhamamura/via"
|
||||||
// "github.com/go-via/via-plugin-picocss/picocss"
|
|
||||||
. "github.com/ryanhamamura/via/h"
|
. "github.com/ryanhamamura/via/h"
|
||||||
)
|
)
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
v := via.New()
|
v := via.New()
|
||||||
|
|
||||||
v.Config(via.Options{
|
|
||||||
// Plugins: []via.Plugin{picocss.Default},
|
|
||||||
})
|
|
||||||
|
|
||||||
v.Page("/counters/{counter_id}/{start_at_step}", func(c *via.Context) {
|
v.Page("/counters/{counter_id}/{start_at_step}", func(c *via.Context) {
|
||||||
|
|
||||||
counterID := c.GetPathParam("counter_id")
|
counterID := c.GetPathParam("counter_id")
|
||||||
|
|||||||
Binary file not shown.
@@ -2,7 +2,6 @@ package main
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"github.com/ryanhamamura/via"
|
"github.com/ryanhamamura/via"
|
||||||
// "github.com/go-via/via-plugin-picocss/picocss"
|
|
||||||
"github.com/ryanhamamura/via/h"
|
"github.com/ryanhamamura/via/h"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -13,11 +12,6 @@ func main() {
|
|||||||
|
|
||||||
v.Config(via.Options{
|
v.Config(via.Options{
|
||||||
DocumentTitle: "Via Counter",
|
DocumentTitle: "Via Counter",
|
||||||
// Plugin is placed here. Use picocss.WithOptions(pococss.Options) to add the plugin
|
|
||||||
// with a different color theme or to enable a classes for a wide range of colors.
|
|
||||||
// Plugins: []via.Plugin{
|
|
||||||
// picocss.Default,
|
|
||||||
// },
|
|
||||||
})
|
})
|
||||||
|
|
||||||
v.Page("/", func(c *via.Context) {
|
v.Page("/", func(c *via.Context) {
|
||||||
|
|||||||
@@ -6,7 +6,6 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ryanhamamura/via"
|
"github.com/ryanhamamura/via"
|
||||||
// "github.com/go-via/via-plugin-picocss/picocss"
|
|
||||||
"github.com/ryanhamamura/via/h"
|
"github.com/ryanhamamura/via/h"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -16,9 +15,6 @@ func main() {
|
|||||||
v.Config(via.Options{
|
v.Config(via.Options{
|
||||||
LogLevel: via.LogLevelDebug,
|
LogLevel: via.LogLevelDebug,
|
||||||
DevMode: true,
|
DevMode: true,
|
||||||
Plugins: []via.Plugin{
|
|
||||||
// picocss.Default,
|
|
||||||
},
|
|
||||||
})
|
})
|
||||||
|
|
||||||
v.AppendToHead(
|
v.AppendToHead(
|
||||||
|
|||||||
@@ -6,9 +6,9 @@ import (
|
|||||||
"log"
|
"log"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
_ "github.com/mattn/go-sqlite3"
|
||||||
"github.com/ryanhamamura/via"
|
"github.com/ryanhamamura/via"
|
||||||
"github.com/ryanhamamura/via/h"
|
"github.com/ryanhamamura/via/h"
|
||||||
_ "github.com/mattn/go-sqlite3"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
type DataSource interface {
|
type DataSource interface {
|
||||||
@@ -22,6 +22,9 @@ type ShakeDB struct {
|
|||||||
findByTextStmt *sql.Stmt
|
findByTextStmt *sql.Stmt
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Prepare opens shake.db, a ~22 MB SQLite database of Shakespeare's works.
|
||||||
|
// Download from https://github.com/nicholasgasior/gopher-fizzbuzz/raw/master/shake.db
|
||||||
|
// and place it in this directory before running.
|
||||||
func (shakeDB *ShakeDB) Prepare() {
|
func (shakeDB *ShakeDB) Prepare() {
|
||||||
db, err := sql.Open("sqlite3", "shake.db")
|
db, err := sql.Open("sqlite3", "shake.db")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
Binary file not shown.
@@ -1,7 +1,6 @@
|
|||||||
package via
|
package via
|
||||||
|
|
||||||
import (
|
import (
|
||||||
// "net/http/httptest"
|
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"github.com/ryanhamamura/via/h"
|
"github.com/ryanhamamura/via/h"
|
||||||
|
|||||||
158
via.go
158
via.go
@@ -58,7 +58,6 @@ type V struct {
|
|||||||
datastarPath string
|
datastarPath string
|
||||||
datastarContent []byte
|
datastarContent []byte
|
||||||
datastarOnce sync.Once
|
datastarOnce sync.Once
|
||||||
reaperStop chan struct{}
|
|
||||||
middleware []Middleware
|
middleware []Middleware
|
||||||
layout func(func() h.H) h.H
|
layout func(func() h.H) h.H
|
||||||
}
|
}
|
||||||
@@ -139,12 +138,6 @@ func (v *V) Config(cfg Options) {
|
|||||||
v.defaultNATS = nil
|
v.defaultNATS = nil
|
||||||
v.pubsub = cfg.PubSub
|
v.pubsub = cfg.PubSub
|
||||||
}
|
}
|
||||||
if cfg.ContextSuspendAfter != 0 {
|
|
||||||
v.cfg.ContextSuspendAfter = cfg.ContextSuspendAfter
|
|
||||||
}
|
|
||||||
if cfg.ContextTTL != 0 {
|
|
||||||
v.cfg.ContextTTL = cfg.ContextTTL
|
|
||||||
}
|
|
||||||
if cfg.Streams != nil {
|
if cfg.Streams != nil {
|
||||||
v.cfg.Streams = cfg.Streams
|
v.cfg.Streams = cfg.Streams
|
||||||
}
|
}
|
||||||
@@ -292,75 +285,6 @@ func (v *V) getCtx(id string) (*Context, error) {
|
|||||||
return nil, fmt.Errorf("ctx '%s' not found", id)
|
return nil, fmt.Errorf("ctx '%s' not found", id)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (v *V) startReaper() {
|
|
||||||
ttl := v.cfg.ContextTTL
|
|
||||||
if ttl < 0 {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
if ttl == 0 {
|
|
||||||
ttl = time.Hour
|
|
||||||
}
|
|
||||||
suspendAfter := v.cfg.ContextSuspendAfter
|
|
||||||
if suspendAfter == 0 {
|
|
||||||
suspendAfter = 15 * time.Minute
|
|
||||||
}
|
|
||||||
if suspendAfter > ttl {
|
|
||||||
suspendAfter = ttl
|
|
||||||
}
|
|
||||||
interval := suspendAfter / 3
|
|
||||||
if interval < 5*time.Second {
|
|
||||||
interval = 5 * time.Second
|
|
||||||
}
|
|
||||||
v.reaperStop = make(chan struct{})
|
|
||||||
go func() {
|
|
||||||
ticker := time.NewTicker(interval)
|
|
||||||
defer ticker.Stop()
|
|
||||||
for {
|
|
||||||
select {
|
|
||||||
case <-v.reaperStop:
|
|
||||||
return
|
|
||||||
case <-ticker.C:
|
|
||||||
v.reapOrphanedContexts(suspendAfter, ttl)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
}
|
|
||||||
|
|
||||||
func (v *V) reapOrphanedContexts(suspendAfter, ttl time.Duration) {
|
|
||||||
now := time.Now()
|
|
||||||
v.contextRegistryMutex.RLock()
|
|
||||||
var toSuspend, toReap []*Context
|
|
||||||
for _, c := range v.contextRegistry {
|
|
||||||
if c.sseConnected.Load() {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
// Use the most recent liveness signal
|
|
||||||
lastAlive := c.createdAt
|
|
||||||
if dc := c.sseDisconnectedAt.Load(); dc != nil && dc.After(lastAlive) {
|
|
||||||
lastAlive = *dc
|
|
||||||
}
|
|
||||||
if seen := c.lastSeenAt.Load(); seen != nil && seen.After(lastAlive) {
|
|
||||||
lastAlive = *seen
|
|
||||||
}
|
|
||||||
silentFor := now.Sub(lastAlive)
|
|
||||||
if silentFor > ttl {
|
|
||||||
toReap = append(toReap, c)
|
|
||||||
} else if silentFor > suspendAfter && !c.suspended.Load() {
|
|
||||||
toSuspend = append(toSuspend, c)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
v.contextRegistryMutex.RUnlock()
|
|
||||||
|
|
||||||
for _, c := range toSuspend {
|
|
||||||
v.logInfo(c, "suspending context (no SSE connection after %s)", suspendAfter)
|
|
||||||
c.suspend()
|
|
||||||
}
|
|
||||||
for _, c := range toReap {
|
|
||||||
v.logInfo(c, "reaping orphaned context (no SSE connection after %s)", ttl)
|
|
||||||
v.cleanupCtx(c)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Start starts the Via HTTP server and blocks until a SIGINT or SIGTERM
|
// Start starts the Via HTTP server and blocks until a SIGINT or SIGTERM
|
||||||
// signal is received, then performs a graceful shutdown.
|
// signal is received, then performs a graceful shutdown.
|
||||||
func (v *V) Start() {
|
func (v *V) Start() {
|
||||||
@@ -389,8 +313,6 @@ func (v *V) Start() {
|
|||||||
Handler: handler,
|
Handler: handler,
|
||||||
}
|
}
|
||||||
|
|
||||||
v.startReaper()
|
|
||||||
|
|
||||||
errCh := make(chan error, 1)
|
errCh := make(chan error, 1)
|
||||||
go func() {
|
go func() {
|
||||||
errCh <- v.server.ListenAndServe()
|
errCh <- v.server.ListenAndServe()
|
||||||
@@ -417,9 +339,6 @@ func (v *V) Start() {
|
|||||||
// Shutdown gracefully shuts down the server and all contexts.
|
// Shutdown gracefully shuts down the server and all contexts.
|
||||||
// Safe for programmatic or test use.
|
// Safe for programmatic or test use.
|
||||||
func (v *V) Shutdown() {
|
func (v *V) Shutdown() {
|
||||||
if v.reaperStop != nil {
|
|
||||||
close(v.reaperStop)
|
|
||||||
}
|
|
||||||
v.logInfo(nil, "draining all contexts")
|
v.logInfo(nil, "draining all contexts")
|
||||||
v.drainAllContexts()
|
v.drainAllContexts()
|
||||||
|
|
||||||
@@ -520,36 +439,39 @@ func (v *V) ensureDatastarHandler() {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func loadDevModeMap(path string) map[string]string {
|
||||||
|
m := make(map[string]string)
|
||||||
|
file, err := os.Open(path)
|
||||||
|
if err != nil {
|
||||||
|
return m
|
||||||
|
}
|
||||||
|
defer file.Close()
|
||||||
|
json.NewDecoder(file).Decode(&m)
|
||||||
|
return m
|
||||||
|
}
|
||||||
|
|
||||||
|
func saveDevModeMap(path string, m map[string]string) error {
|
||||||
|
file, err := os.Create(path)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
defer file.Close()
|
||||||
|
return json.NewEncoder(file).Encode(m)
|
||||||
|
}
|
||||||
|
|
||||||
func (v *V) devModePersist(c *Context) {
|
func (v *V) devModePersist(c *Context) {
|
||||||
p := filepath.Join(".via", "devmode", "ctx.json")
|
p := filepath.Join(".via", "devmode", "ctx.json")
|
||||||
if err := os.MkdirAll(filepath.Dir(p), 0755); err != nil {
|
if err := os.MkdirAll(filepath.Dir(p), 0755); err != nil {
|
||||||
v.logFatal("failed to create directory for devmode files: %v", err)
|
v.logFatal("failed to create directory for devmode files: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// load persisted list from file, or empty list if file not found
|
ctxRegMap := loadDevModeMap(p)
|
||||||
file, err := os.Open(p)
|
|
||||||
ctxRegMap := make(map[string]string)
|
|
||||||
if err == nil {
|
|
||||||
json.NewDecoder(file).Decode(&ctxRegMap)
|
|
||||||
}
|
|
||||||
file.Close()
|
|
||||||
|
|
||||||
// add ctx to persisted list
|
|
||||||
if _, ok := ctxRegMap[c.id]; !ok {
|
if _, ok := ctxRegMap[c.id]; !ok {
|
||||||
ctxRegMap[c.id] = c.route
|
ctxRegMap[c.id] = c.route
|
||||||
}
|
}
|
||||||
|
|
||||||
// write persisted list to file
|
if err := saveDevModeMap(p, ctxRegMap); err != nil {
|
||||||
file, err = os.Create(p)
|
v.logErr(c, "devmode failed to persist ctx: %v", err)
|
||||||
if err != nil {
|
|
||||||
v.logErr(c, "devmode failed to percist ctx: %v", err)
|
|
||||||
|
|
||||||
}
|
|
||||||
defer file.Close()
|
|
||||||
|
|
||||||
encoder := json.NewEncoder(file)
|
|
||||||
if err := encoder.Encode(ctxRegMap); err != nil {
|
|
||||||
v.logErr(c, "devmode failed to persist ctx")
|
|
||||||
}
|
}
|
||||||
v.logDebug(c, "devmode persisted ctx to file")
|
v.logDebug(c, "devmode persisted ctx to file")
|
||||||
}
|
}
|
||||||
@@ -557,27 +479,11 @@ func (v *V) devModePersist(c *Context) {
|
|||||||
func (v *V) devModeRemovePersisted(c *Context) {
|
func (v *V) devModeRemovePersisted(c *Context) {
|
||||||
p := filepath.Join(".via", "devmode", "ctx.json")
|
p := filepath.Join(".via", "devmode", "ctx.json")
|
||||||
|
|
||||||
// load persisted list from file, or empty list if file not found
|
ctxRegMap := loadDevModeMap(p)
|
||||||
file, err := os.Open(p)
|
|
||||||
ctxRegMap := make(map[string]string)
|
|
||||||
if err == nil {
|
|
||||||
json.NewDecoder(file).Decode(&ctxRegMap)
|
|
||||||
}
|
|
||||||
file.Close()
|
|
||||||
|
|
||||||
delete(ctxRegMap, c.id)
|
delete(ctxRegMap, c.id)
|
||||||
|
|
||||||
// write persisted list to file
|
if err := saveDevModeMap(p, ctxRegMap); err != nil {
|
||||||
file, err = os.Create(p)
|
v.logErr(c, "devmode failed to remove persisted ctx: %v", err)
|
||||||
if err != nil {
|
|
||||||
v.logErr(c, "devmode failed to remove percisted ctx: %v", err)
|
|
||||||
|
|
||||||
}
|
|
||||||
defer file.Close()
|
|
||||||
|
|
||||||
encoder := json.NewEncoder(file)
|
|
||||||
if err := encoder.Encode(ctxRegMap); err != nil {
|
|
||||||
v.logErr(c, "devmode failed to remove persisted ctx")
|
|
||||||
}
|
}
|
||||||
v.logDebug(c, "devmode removed persisted ctx from file")
|
v.logDebug(c, "devmode removed persisted ctx from file")
|
||||||
}
|
}
|
||||||
@@ -667,8 +573,6 @@ func New() *V {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
c.reqCtx = r.Context()
|
c.reqCtx = r.Context()
|
||||||
now := time.Now()
|
|
||||||
c.lastSeenAt.Store(&now)
|
|
||||||
|
|
||||||
sse := datastar.NewSSE(w, r, datastar.WithCompression(datastar.WithBrotli(datastar.WithBrotliLevel(5))))
|
sse := datastar.NewSSE(w, r, datastar.WithCompression(datastar.WithBrotli(datastar.WithBrotliLevel(5))))
|
||||||
|
|
||||||
@@ -690,16 +594,6 @@ func New() *V {
|
|||||||
c.sseDisconnectedAt.Store(nil)
|
c.sseDisconnectedAt.Store(nil)
|
||||||
v.logDebug(c, "SSE connection established")
|
v.logDebug(c, "SSE connection established")
|
||||||
|
|
||||||
if c.suspended.Load() {
|
|
||||||
c.navMu.Lock()
|
|
||||||
c.suspended.Store(false)
|
|
||||||
if initFn := v.pageRegistry[c.route]; initFn != nil {
|
|
||||||
v.logInfo(c, "resuming suspended context")
|
|
||||||
initFn(c)
|
|
||||||
}
|
|
||||||
c.navMu.Unlock()
|
|
||||||
}
|
|
||||||
|
|
||||||
go c.Sync()
|
go c.Sync()
|
||||||
|
|
||||||
keepalive := time.NewTicker(30 * time.Second)
|
keepalive := time.NewTicker(30 * time.Second)
|
||||||
|
|||||||
157
via_test.go
157
via_test.go
@@ -400,95 +400,6 @@ func TestPage_PanicsOnNoView(t *testing.T) {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestReaperCleansOrphanedContexts(t *testing.T) {
|
|
||||||
v := New()
|
|
||||||
c := newContext("orphan-1", "/", v)
|
|
||||||
c.createdAt = time.Now().Add(-time.Minute) // created 1 min ago
|
|
||||||
v.registerCtx(c)
|
|
||||||
|
|
||||||
_, err := v.getCtx("orphan-1")
|
|
||||||
assert.NoError(t, err)
|
|
||||||
|
|
||||||
v.reapOrphanedContexts(5*time.Second, 10*time.Second)
|
|
||||||
|
|
||||||
_, err = v.getCtx("orphan-1")
|
|
||||||
assert.Error(t, err, "orphaned context should have been reaped")
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestReaperSuspendsContext(t *testing.T) {
|
|
||||||
v := New()
|
|
||||||
c := newContext("suspend-1", "/", v)
|
|
||||||
c.createdAt = time.Now().Add(-30 * time.Minute)
|
|
||||||
dc := time.Now().Add(-20 * time.Minute)
|
|
||||||
c.sseDisconnectedAt.Store(&dc)
|
|
||||||
v.registerCtx(c)
|
|
||||||
|
|
||||||
v.reapOrphanedContexts(15*time.Minute, time.Hour)
|
|
||||||
|
|
||||||
got, err := v.getCtx("suspend-1")
|
|
||||||
assert.NoError(t, err, "suspended context should still be in registry")
|
|
||||||
assert.True(t, got.suspended.Load(), "context should be marked suspended")
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestReaperReapsAfterTTL(t *testing.T) {
|
|
||||||
v := New()
|
|
||||||
c := newContext("reap-1", "/", v)
|
|
||||||
c.createdAt = time.Now().Add(-3 * time.Hour)
|
|
||||||
dc := time.Now().Add(-2 * time.Hour)
|
|
||||||
c.sseDisconnectedAt.Store(&dc)
|
|
||||||
c.suspended.Store(true)
|
|
||||||
v.registerCtx(c)
|
|
||||||
|
|
||||||
v.reapOrphanedContexts(15*time.Minute, time.Hour)
|
|
||||||
|
|
||||||
_, err := v.getCtx("reap-1")
|
|
||||||
assert.Error(t, err, "context past TTL should have been reaped")
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestReaperIgnoresAlreadySuspended(t *testing.T) {
|
|
||||||
v := New()
|
|
||||||
c := newContext("already-sus-1", "/", v)
|
|
||||||
c.createdAt = time.Now().Add(-30 * time.Minute)
|
|
||||||
dc := time.Now().Add(-20 * time.Minute)
|
|
||||||
c.sseDisconnectedAt.Store(&dc)
|
|
||||||
c.suspended.Store(true)
|
|
||||||
// give it a fresh pageStopChan so we can verify it's not re-closed
|
|
||||||
c.pageStopChan = make(chan struct{})
|
|
||||||
v.registerCtx(c)
|
|
||||||
|
|
||||||
v.reapOrphanedContexts(15*time.Minute, time.Hour)
|
|
||||||
|
|
||||||
got, err := v.getCtx("already-sus-1")
|
|
||||||
assert.NoError(t, err, "already-suspended context within TTL should survive")
|
|
||||||
assert.True(t, got.suspended.Load())
|
|
||||||
// pageStopChan should still be open (not re-suspended)
|
|
||||||
select {
|
|
||||||
case <-got.pageStopChan:
|
|
||||||
t.Fatal("pageStopChan was closed — context was re-suspended")
|
|
||||||
default:
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestReaperIgnoresConnectedContexts(t *testing.T) {
|
|
||||||
v := New()
|
|
||||||
c := newContext("connected-1", "/", v)
|
|
||||||
c.createdAt = time.Now().Add(-time.Minute)
|
|
||||||
c.sseConnected.Store(true)
|
|
||||||
v.registerCtx(c)
|
|
||||||
|
|
||||||
v.reapOrphanedContexts(5*time.Second, 10*time.Second)
|
|
||||||
|
|
||||||
_, err := v.getCtx("connected-1")
|
|
||||||
assert.NoError(t, err, "connected context should survive reaping")
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestReaperDisabledWithNegativeTTL(t *testing.T) {
|
|
||||||
v := New()
|
|
||||||
v.cfg.ContextTTL = -1
|
|
||||||
v.startReaper()
|
|
||||||
assert.Nil(t, v.reaperStop, "reaper should not start with negative TTL")
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestCleanupCtxIdempotent(t *testing.T) {
|
func TestCleanupCtxIdempotent(t *testing.T) {
|
||||||
v := New()
|
v := New()
|
||||||
c := newContext("idempotent-1", "/", v)
|
c := newContext("idempotent-1", "/", v)
|
||||||
@@ -503,74 +414,6 @@ func TestCleanupCtxIdempotent(t *testing.T) {
|
|||||||
assert.Error(t, err, "context should be removed after cleanup")
|
assert.Error(t, err, "context should be removed after cleanup")
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestReaperRespectsLastSeenAt(t *testing.T) {
|
|
||||||
v := New()
|
|
||||||
c := newContext("seen-1", "/", v)
|
|
||||||
c.createdAt = time.Now().Add(-30 * time.Minute)
|
|
||||||
// Disconnected 20 min ago, but client retried (lastSeenAt) 2 min ago
|
|
||||||
dc := time.Now().Add(-20 * time.Minute)
|
|
||||||
c.sseDisconnectedAt.Store(&dc)
|
|
||||||
seen := time.Now().Add(-2 * time.Minute)
|
|
||||||
c.lastSeenAt.Store(&seen)
|
|
||||||
v.registerCtx(c)
|
|
||||||
|
|
||||||
v.reapOrphanedContexts(15*time.Minute, time.Hour)
|
|
||||||
|
|
||||||
_, err := v.getCtx("seen-1")
|
|
||||||
assert.NoError(t, err, "context with recent lastSeenAt should survive suspend threshold")
|
|
||||||
assert.False(t, c.suspended.Load(), "context should not be suspended")
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestReaperFallsBackWithoutLastSeenAt(t *testing.T) {
|
|
||||||
v := New()
|
|
||||||
c := newContext("noseen-1", "/", v)
|
|
||||||
c.createdAt = time.Now().Add(-30 * time.Minute)
|
|
||||||
dc := time.Now().Add(-20 * time.Minute)
|
|
||||||
c.sseDisconnectedAt.Store(&dc)
|
|
||||||
// no lastSeenAt set — should fall back to sseDisconnectedAt
|
|
||||||
v.registerCtx(c)
|
|
||||||
|
|
||||||
v.reapOrphanedContexts(15*time.Minute, time.Hour)
|
|
||||||
|
|
||||||
got, err := v.getCtx("noseen-1")
|
|
||||||
assert.NoError(t, err, "context should still be in registry (suspended, not reaped)")
|
|
||||||
assert.True(t, got.suspended.Load(), "context should be suspended using sseDisconnectedAt fallback")
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestReaperReapsWithStaleLastSeenAt(t *testing.T) {
|
|
||||||
v := New()
|
|
||||||
c := newContext("stale-seen-1", "/", v)
|
|
||||||
c.createdAt = time.Now().Add(-3 * time.Hour)
|
|
||||||
dc := time.Now().Add(-2 * time.Hour)
|
|
||||||
c.sseDisconnectedAt.Store(&dc)
|
|
||||||
// lastSeenAt is also old — beyond TTL
|
|
||||||
seen := time.Now().Add(-90 * time.Minute)
|
|
||||||
c.lastSeenAt.Store(&seen)
|
|
||||||
c.suspended.Store(true)
|
|
||||||
v.registerCtx(c)
|
|
||||||
|
|
||||||
v.reapOrphanedContexts(15*time.Minute, time.Hour)
|
|
||||||
|
|
||||||
_, err := v.getCtx("stale-seen-1")
|
|
||||||
assert.Error(t, err, "context with stale lastSeenAt beyond TTL should be reaped")
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestLastSeenAtUpdatedOnSSEConnect(t *testing.T) {
|
|
||||||
v := New()
|
|
||||||
c := newContext("seen-sse-1", "/", v)
|
|
||||||
v.registerCtx(c)
|
|
||||||
|
|
||||||
assert.Nil(t, c.lastSeenAt.Load(), "lastSeenAt should be nil before SSE connect")
|
|
||||||
|
|
||||||
// Simulate what the SSE handler does after getCtx
|
|
||||||
now := time.Now()
|
|
||||||
c.lastSeenAt.Store(&now)
|
|
||||||
|
|
||||||
got := c.lastSeenAt.Load()
|
|
||||||
assert.NotNil(t, got, "lastSeenAt should be set after SSE connect")
|
|
||||||
assert.WithinDuration(t, now, *got, time.Second)
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestDevModeRemovePersistedFix(t *testing.T) {
|
func TestDevModeRemovePersistedFix(t *testing.T) {
|
||||||
v := New()
|
v := New()
|
||||||
v.cfg.DevMode = true
|
v.cfg.DevMode = true
|
||||||
|
|||||||
Reference in New Issue
Block a user