HH-807: harden runtime coordination (#21)

This commit is contained in:
2026-08-31 19:37:17 +08:00
parent 1d9fd9f0e0
commit 024bf14448
7 changed files with 937 additions and 45 deletions
+287 -26
View File
@@ -10,7 +10,6 @@ import (
"net/http"
"regexp"
"strings"
"sync"
"time"
"git.ipao.vip/rogee/creator-hub/internal/hub"
@@ -19,6 +18,7 @@ import (
// hubStore 是控制面编排所需的存储能力;生产实现为 *hub.Store,测试使用内存桩。
type hubStore interface {
LockResources(ctx context.Context, aliases, exitIDs, imageVersions []string) (func(), error)
CreateGateway(ctx context.Context, name, endpoint, token string) (hub.Gateway, error)
ListGateways(ctx context.Context) ([]hub.Gateway, error)
GetGateway(ctx context.Context, name string) (hub.Gateway, error)
@@ -50,6 +50,7 @@ type hubStore interface {
}
type runtimeStopStore interface {
LockResources(ctx context.Context, aliases, exitIDs, imageVersions []string) (func(), error)
GetEnvironmentContextForAccount(ctx context.Context, accountID string) (hub.EnvironmentContext, error)
GetGateway(ctx context.Context, name string) (hub.Gateway, error)
ReleaseRuntime(ctx context.Context, environment hub.EnvironmentContext) error
@@ -369,31 +370,37 @@ func releaseRuntime(ctx context.Context, store runtimeStopStore, environment hub
return err
}
func releaseRuntimeWithReconcileAudit(ctx context.Context, store hubStore, environment hub.EnvironmentContext) error {
action := actionForEnvironment("reconcile", environment)
if err := store.AppendEnvironmentAction(ctx, "environment_action_requested", action); err != nil {
return err
}
releaseErr := releaseRuntime(ctx, store, environment)
action.Outcome, action.ReasonCode = "succeeded", "runtime_released"
if releaseErr != nil {
action.Outcome, action.ReasonCode = "failed", "runtime_release_failed"
}
if auditErr := store.AppendEnvironmentAction(ctx, "environment_action_finished", action); auditErr != nil {
return errors.Join(releaseErr, auditErr)
}
return releaseErr
}
func registerHub(app *fiber.App, store hubStore) {
registerHubWithNetwork(app, store, defaultNetworkExitProbe(), resolveExitCredential)
}
var runtimeOperations sync.Mutex
func registerHubWithNetwork(app *fiber.App, store hubStore, probe networkExitProbe, resolve func(hub.NetworkExitAccess) (string, error)) {
// ponytail: one control-plane instance is serialized globally; use keyed/distributed locks if replicas or throughput require it.
serialized := func(handler fiber.Handler) fiber.Handler {
return func(c fiber.Ctx) error {
runtimeOperations.Lock()
defer runtimeOperations.Unlock()
return handler(c)
}
}
app.Get("/api/browsers", serialized(listBrowsers(store, probe, resolve)))
app.Get("/api/browsers", listBrowsers(store, probe, resolve))
app.Get("/api/browsers/:alias", getBrowser(store))
app.Post("/api/browsers", serialized(createBrowser(store, probe, resolve)))
app.Post("/api/browsers/:alias/:action", serialized(browserAction(store, probe, resolve)))
app.Delete("/api/browsers/:alias", serialized(deleteBrowser(store)))
app.Post("/api/browsers", lockBrowserCreate(store, createBrowser(store, probe, resolve)))
app.Post("/api/browsers/:alias/:action", lockBrowserAlias(store, browserAction(store, probe, resolve)))
app.Delete("/api/browsers/:alias", lockBrowserAlias(store, deleteBrowser(store)))
app.Get("/api/network-exits", listNetworkExits(store))
app.Get("/api/network-exits/:id", getNetworkExit(store))
app.Post("/api/network-exits", serialized(createNetworkExit(store)))
app.Post("/api/network-exits/:id/check", serialized(checkNetworkExit(store, probe)))
app.Post("/api/network-exits/:id/disable", serialized(disableNetworkExit(store, probe, resolve)))
app.Post("/api/network-exits", createNetworkExit(store))
app.Post("/api/network-exits/:id/check", checkNetworkExit(store, probe))
app.Post("/api/network-exits/:id/disable", disableNetworkExit(store, probe, resolve))
app.Get("/api/browser-images", func(c fiber.Ctx) error {
images, err := store.ListImages(c.Context(), false)
@@ -423,7 +430,7 @@ func registerHubWithNetwork(app *fiber.App, store hubStore, probe networkExitPro
"version": input.Version, "image_ref": input.ImageRef, "note": input.Note, "enabled": enabled,
})
})
app.Put("/api/browser-images/:version", serialized(func(c fiber.Ctx) error {
app.Put("/api/browser-images/:version", lockBrowserImage(store, func(c fiber.Ctx) error {
input := struct {
ImageRef string `json:"image_ref"`
Note string `json:"note"`
@@ -441,7 +448,7 @@ func registerHubWithNetwork(app *fiber.App, store hubStore, probe networkExitPro
}
return c.SendStatus(fiber.StatusNoContent)
}))
app.Delete("/api/browser-images/:version", serialized(func(c fiber.Ctx) error {
app.Delete("/api/browser-images/:version", lockBrowserImage(store, func(c fiber.Ctx) error {
if err := store.DeleteImage(c.Context(), c.Params("version")); err != nil {
return hubError(c, err)
}
@@ -488,6 +495,236 @@ func getBrowser(store hubStore) fiber.Handler {
}
}
func lockBrowserAlias(store hubStore, handler fiber.Handler) fiber.Handler {
return func(c fiber.Ctx) error {
requestedExitID, imageVersion := "", ""
if c.Params("action") == "rebind" {
var input struct {
NetworkExitID string `json:"network_exit_id"`
}
if json.Unmarshal(c.Body(), &input) == nil {
if hub.ValidNetworkExitID(input.NetworkExitID) {
requestedExitID = input.NetworkExitID
}
}
}
if c.Params("action") == "upgrade" {
var input struct {
Version string `json:"version"`
}
if json.Unmarshal(c.Body(), &input) == nil {
if hub.ValidImageVersion(input.Version) {
imageVersion = input.Version
}
}
}
alias := c.Params("alias")
unlock, err := lockAliasResources(c.Context(), store, alias, requestedExitID, imageVersion)
if err != nil {
return hubError(c, err)
}
defer unlock()
return handler(c)
}
}
func lockBrowserImage(store hubStore, handler fiber.Handler) fiber.Handler {
return func(c fiber.Ctx) error {
unlock, err := store.LockResources(c.Context(), nil, nil, []string{c.Params("version")})
if err != nil {
return hubError(c, err)
}
defer unlock()
return handler(c)
}
}
func lockBrowserCreate(store hubStore, handler fiber.Handler) fiber.Handler {
return func(c fiber.Ctx) error {
var input struct {
Alias string `json:"alias"`
AccountID string `json:"account_id"`
NetworkExitID string `json:"network_exit_id"`
ImageVersion string `json:"image_version"`
}
if json.Unmarshal(c.Body(), &input) != nil || input.Alias == "" {
return handler(c)
}
exitIDs, imageVersions := []string(nil), []string(nil)
if hub.ValidNetworkExitID(input.NetworkExitID) {
exitIDs = []string{input.NetworkExitID}
}
if hub.ValidImageVersion(input.ImageVersion) {
imageVersions = []string{input.ImageVersion}
}
unlock, err := store.LockResources(c.Context(), nonEmpty(input.Alias, input.AccountID), exitIDs, imageVersions)
if err != nil {
return hubError(c, err)
}
defer unlock()
return handler(c)
}
}
func lockEnvironmentResources(ctx context.Context, store hubStore, envs []hub.Env) (func(), error) {
aliases := make([]string, 0, len(envs))
seen := make(map[string]bool, len(envs))
for _, env := range envs {
if !seen[env.Alias] {
seen[env.Alias] = true
aliases = append(aliases, env.Alias)
}
}
for {
before, err := environmentResourceMap(ctx, store, aliases)
if err != nil {
return nil, err
}
exitIDs, imageVersions := environmentResourceValues(before)
unlock, err := store.LockResources(ctx, aliases, exitIDs, imageVersions)
if err != nil {
return nil, err
}
after, err := environmentResourceMap(ctx, store, aliases)
if err != nil {
unlock()
return nil, err
}
if equalEnvironmentResourceMaps(before, after) {
return unlock, nil
}
unlock()
}
}
type environmentResource struct {
alias string
exitID string
imageVersion string
}
type resourceLockStore interface {
LockResources(ctx context.Context, aliases, exitIDs, imageVersions []string) (func(), error)
GetEnvironmentContext(ctx context.Context, alias string) (hub.EnvironmentContext, error)
}
func lockAliasResources(ctx context.Context, store resourceLockStore, alias, requestedExitID, requestedImageVersion string) (func(), error) {
for {
before, beforeFound, err := environmentResources(ctx, store, alias)
if err != nil {
return nil, err
}
unlock, err := store.LockResources(ctx, []string{alias}, nonEmpty(before.exitID, requestedExitID), nonEmpty(before.imageVersion, requestedImageVersion))
if err != nil {
return nil, err
}
after, afterFound, err := environmentResources(ctx, store, alias)
if err != nil {
unlock()
return nil, err
}
if beforeFound == afterFound && before == after {
return unlock, nil
}
unlock()
}
}
func lockAccountResources(ctx context.Context, store runtimeStopStore, accountID string) (func(), error) {
if store == nil {
return func() {}, nil
}
for {
before, beforeFound, err := accountEnvironmentResources(ctx, store, accountID)
if err != nil {
return nil, err
}
unlock, err := store.LockResources(ctx, nonEmpty(accountID, before.alias), nonEmpty(before.exitID), nonEmpty(before.imageVersion))
if err != nil {
return nil, err
}
after, afterFound, err := accountEnvironmentResources(ctx, store, accountID)
if err != nil {
unlock()
return nil, err
}
if beforeFound == afterFound && before == after {
return unlock, nil
}
unlock()
}
}
func accountEnvironmentResources(ctx context.Context, store runtimeStopStore, accountID string) (environmentResource, bool, error) {
environment, err := store.GetEnvironmentContextForAccount(ctx, accountID)
if errors.Is(err, hub.ErrNotFound) {
return environmentResource{}, false, nil
}
if err != nil {
return environmentResource{}, false, err
}
return environmentResource{alias: environment.Alias, exitID: environment.Exit.ID, imageVersion: environment.ImageVersion}, true, nil
}
func environmentResources(ctx context.Context, store resourceLockStore, alias string) (environmentResource, bool, error) {
environment, err := store.GetEnvironmentContext(ctx, alias)
if errors.Is(err, hub.ErrNotFound) {
return environmentResource{}, false, nil
}
if err != nil {
return environmentResource{}, false, err
}
return environmentResource{alias: environment.Alias, exitID: environment.Exit.ID, imageVersion: environment.ImageVersion}, true, nil
}
func environmentResourceMap(ctx context.Context, store hubStore, aliases []string) (map[string]environmentResource, error) {
result := make(map[string]environmentResource, len(aliases))
for _, alias := range aliases {
resources, found, err := environmentResources(ctx, store, alias)
if err != nil {
return nil, err
}
if found {
result[alias] = resources
}
}
return result, nil
}
func nonEmpty(values ...string) []string {
result := make([]string, 0, len(values))
seen := make(map[string]bool, len(values))
for _, value := range values {
if value != "" && !seen[value] {
seen[value] = true
result = append(result, value)
}
}
return result
}
func environmentResourceValues(values map[string]environmentResource) ([]string, []string) {
exitIDs, imageVersions := make([]string, 0, len(values)), make([]string, 0, len(values))
for _, value := range values {
exitIDs = append(exitIDs, value.exitID)
imageVersions = append(imageVersions, value.imageVersion)
}
return nonEmpty(exitIDs...), nonEmpty(imageVersions...)
}
func equalEnvironmentResourceMaps(left, right map[string]environmentResource) bool {
if len(left) != len(right) {
return false
}
for key, value := range left {
other, ok := right[key]
if !ok || other != value {
return false
}
}
return true
}
func listNetworkExits(store hubStore) fiber.Handler {
return func(c fiber.Ctx) error {
exits, err := store.ListNetworkExits(c.Context())
@@ -536,6 +773,11 @@ func createNetworkExit(store hubStore) fiber.Handler {
func checkNetworkExit(store hubStore, probe networkExitProbe) fiber.Handler {
return func(c fiber.Ctx) error {
unlock, err := store.LockResources(c.Context(), nil, []string{c.Params("id")}, nil)
if err != nil {
return hubError(c, err)
}
defer unlock()
access, err := store.GetNetworkExitAccess(c.Context(), c.Params("id"))
if err != nil {
return hubError(c, err)
@@ -557,7 +799,12 @@ func checkNetworkExit(store hubStore, probe networkExitProbe) fiber.Handler {
func disableNetworkExit(store hubStore, probe networkExitProbe, resolve func(hub.NetworkExitAccess) (string, error)) fiber.Handler {
return func(c fiber.Ctx) error {
unlock, err := store.LockResources(c.Context(), nil, []string{c.Params("id")}, nil)
if err != nil {
return hubError(c, err)
}
exit, err := store.DisableNetworkExit(c.Context(), c.Params("id"))
unlock()
if err != nil {
return hubError(c, err)
}
@@ -594,6 +841,11 @@ func listBrowsers(store hubStore, probe networkExitProbe, resolve func(hub.Netwo
if err != nil {
return hubError(c, err)
}
unlock, err := lockEnvironmentResources(c.Context(), store, envs)
if err != nil {
return hubError(c, err)
}
defer unlock()
containers, gatewayRead, gatewayErrors := gatewayContainerSnapshot(c.Context(), store, envs)
if err := reconcileRuntimeSnapshot(c.Context(), store, probe, resolve, envs, containers, gatewayRead, gatewayErrors); err != nil {
return hubError(c, err)
@@ -622,8 +874,6 @@ func listBrowsers(store hubStore, probe networkExitProbe, resolve func(hub.Netwo
}
func reconcileRuntimeLeases(ctx context.Context, store hubStore, probe networkExitProbe, resolve func(hub.NetworkExitAccess) (string, error)) error {
runtimeOperations.Lock()
defer runtimeOperations.Unlock()
return reconcileRuntimeLeasesUnlocked(ctx, store, probe, resolve)
}
@@ -632,6 +882,11 @@ func reconcileRuntimeLeasesUnlocked(ctx context.Context, store hubStore, probe n
if err != nil {
return err
}
unlock, err := lockEnvironmentResources(ctx, store, envs)
if err != nil {
return err
}
defer unlock()
containers, gatewayRead, gatewayErrors := gatewayContainerSnapshot(ctx, store, envs)
return reconcileRuntimeSnapshot(ctx, store, probe, resolve, envs, containers, gatewayRead, gatewayErrors)
}
@@ -708,8 +963,8 @@ func reconcileRuntimeSnapshot(ctx context.Context, store hubStore, probe network
if err := stopEnvironmentRuntime(ctx, store, environment); err != nil {
return err
}
} else if gatewayRead[env.Gateway] {
if err := releaseRuntime(ctx, store, environment); err != nil {
} else if gatewayRead[env.Gateway] && environment.RuntimeInstanceID != "" {
if err := releaseRuntimeWithReconcileAudit(ctx, store, environment); err != nil {
return err
}
}
@@ -752,8 +1007,8 @@ func reconcileRuntimeSnapshot(ctx context.Context, store hubStore, probe network
if restoreErr != nil {
return restoreErr
}
} else if gatewayRead[env.Gateway] {
if err := releaseRuntime(ctx, store, environment); err != nil {
} else if gatewayRead[env.Gateway] && environment.RuntimeInstanceID != "" {
if err := releaseRuntimeWithReconcileAudit(ctx, store, environment); err != nil {
return err
}
}
@@ -1622,6 +1877,12 @@ func rebindBrowser(store hubStore, probe networkExitProbe, resolve func(hub.Netw
action.NetworkExitID, action.BindingVersion, action.RuntimeInstanceID = current.Exit.ID, current.BindingVersion, current.RuntimeInstanceID
return store.AppendEnvironmentAction(c.Context(), "environment_action_finished", action)
}
if !hub.ValidNetworkExitID(input.NetworkExitID) {
if err := finish("failed", "rebind_input_rejected", before); err != nil {
return hubError(c, err)
}
return hubError(c, hub.ErrInvalid)
}
access, reason, err := verifyNetworkExit(c.Context(), store, probe, input.NetworkExitID)
if err != nil {
_ = finish("failed", reason, before)
+440 -1
View File
@@ -11,6 +11,7 @@ import (
"net/http/httptest"
"net/url"
"os"
"sort"
"strings"
"sync"
"testing"
@@ -25,6 +26,8 @@ import (
// memoryStore 是 hubStore 的内存桩,记录写入以便断言编排副作用。
type memoryStore struct {
mu sync.Mutex
locksMu sync.Mutex
locks map[string]*sync.Mutex
gateways map[string]hub.Gateway
images map[string]hub.Image
envs map[string]hub.Env
@@ -53,6 +56,7 @@ func (s *blockingRuntimeStopStore) GetGateway(context.Context, string) (hub.Gate
func newMemoryStore() *memoryStore {
return &memoryStore{
locks: map[string]*sync.Mutex{},
gateways: map[string]hub.Gateway{},
images: map[string]hub.Image{},
envs: map[string]hub.Env{},
@@ -109,6 +113,45 @@ func TestPhaseAReadinessErrorsAreStructured(t *testing.T) {
}
}
func (s *memoryStore) LockResources(_ context.Context, aliases, exitIDs, imageVersions []string) (func(), error) {
keys := make([]string, 0, len(aliases)+len(exitIDs)+len(imageVersions))
for _, alias := range aliases {
keys = append(keys, "environment:"+alias)
}
for _, id := range exitIDs {
keys = append(keys, "network-exit:"+id)
}
for _, version := range imageVersions {
keys = append(keys, "image:"+version)
}
sort.Strings(keys)
unlocks := make([]func(), 0, len(keys))
seen := map[string]bool{}
for _, key := range keys {
if !seen[key] {
seen[key] = true
unlocks = append(unlocks, s.lock(key))
}
}
return func() {
for index := len(unlocks) - 1; index >= 0; index-- {
unlocks[index]()
}
}, nil
}
func (s *memoryStore) lock(key string) func() {
s.locksMu.Lock()
lock := s.locks[key]
if lock == nil {
lock = &sync.Mutex{}
s.locks[key] = lock
}
s.locksMu.Unlock()
lock.Lock()
return lock.Unlock
}
func (s *memoryStore) CreateGateway(_ context.Context, _, _, _ string) (hub.Gateway, error) {
return hub.Gateway{}, nil
}
@@ -556,6 +599,18 @@ type fakeExitProbe struct {
failure string
}
type blockingExitProbe struct {
started chan struct{}
release <-chan struct{}
once sync.Once
}
func (probe *blockingExitProbe) Check(context.Context, hub.NetworkExitAccess) (hub.ExitObservation, string) {
probe.once.Do(func() { close(probe.started) })
<-probe.release
return hub.ExitObservation{PublicIP: "203.0.113.1", Region: "test"}, ""
}
func (probe fakeExitProbe) Check(context.Context, hub.NetworkExitAccess) (hub.ExitObservation, string) {
if probe.observation.PublicIP == "" && probe.failure == "" {
probe.observation = hub.ExitObservation{PublicIP: "203.0.113.1", Region: "test"}
@@ -1583,6 +1638,15 @@ type cleanupCommitUnknownStore struct {
hubStore
}
type failingRuntimeReleaseStore struct {
*hub.Store
err error
}
func (s failingRuntimeReleaseStore) ReleaseRuntime(context.Context, hub.EnvironmentContext) error {
return s.err
}
func setFixtureAccountStatus(t *testing.T, databaseURL, status string) {
t.Helper()
db, err := sql.Open("pgx", databaseURL)
@@ -1609,7 +1673,7 @@ type failContextRefreshStore struct {
func (s *failContextRefreshStore) GetEnvironmentContext(ctx context.Context, alias string) (hub.EnvironmentContext, error) {
s.reads++
if s.reads >= 3 {
if s.reads >= 5 {
return hub.EnvironmentContext{}, errors.New("context refresh unavailable")
}
return s.hubStore.GetEnvironmentContext(ctx, alias)
@@ -3221,6 +3285,295 @@ func TestUpgradeRejectsInvalidVersionWithoutAuditingRawInput(t *testing.T) {
}
}
func TestPostgresRejectsInvalidLifecycleTargetsWithSanitizedAuditPairs(t *testing.T) {
databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL")
if databaseURL == "" {
t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run PostgreSQL integration coverage")
}
const secret = "http://operator:secret@proxy.example"
for _, test := range []struct {
action string
body string
reason string
}{
{action: "upgrade", body: `{"version":"` + secret + `"}`, reason: "upgrade_input_rejected"},
{action: "rebind", body: `{"network_exit_id":"` + secret + `"}`, reason: "rebind_input_rejected"},
} {
t.Run(test.action, func(t *testing.T) {
fixture := newPostgresRebindFixture(t, databaseURL)
app := fiber.New()
registerHubWithNetwork(app, fixture.store, fakeExitProbe{}, func(hub.NetworkExitAccess) (string, error) { return "", nil })
response := do(app, http.MethodPost, "/api/browsers/account-a/"+test.action, test.body)
if response.Code != http.StatusBadRequest {
t.Fatalf("invalid %s returned %d: %s", test.action, response.Code, response.Body.String())
}
var events, operations int
var reason string
if err := fixture.db.QueryRow(`
SELECT count(*), count(DISTINCT operation_id),
coalesce(max(reason_code) FILTER (WHERE event_type = 'environment_action_finished'), '')
FROM audit_event WHERE action = $1`, test.action).Scan(&events, &operations, &reason); err != nil {
t.Fatal(err)
}
if events != 2 || operations != 1 || reason != test.reason {
t.Fatalf("invalid %s audit mismatch: events=%d operations=%d reason=%q", test.action, events, operations, reason)
}
var leaked int
if err := fixture.db.QueryRow(`SELECT count(*) FROM audit_event WHERE row_to_json(audit_event)::text LIKE '%' || $1 || '%'`, secret).Scan(&leaked); err != nil {
t.Fatal(err)
}
if leaked != 0 {
t.Fatalf("invalid %s target leaked into audit", test.action)
}
})
}
}
func TestPostgresReconcileFinishedAuditUsesActivatedRuntime(t *testing.T) {
databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL")
if databaseURL == "" {
t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run PostgreSQL integration coverage")
}
fixture := newPostgresRebindFixture(t, databaseURL)
setFixtureAccountStatus(t, fixture.databaseURL, "active")
var err error
fixture.bound, err = fixture.store.ActivateRuntime(context.Background(), "account-a", "stale-container",
fixture.bound.BindingVersion, fixture.bound.Exit.ID, "network-old")
if err != nil {
t.Fatal(err)
}
fixture.gateway.containers = []containerStatus{{
ID: "stale-container", Alias: "account-a", State: "running", ProxyReady: true,
BindingVersion: fixture.bound.BindingVersion, NetworkExitID: "stale-exit", NetworkID: "network-old",
}}
app := fiber.New()
registerHubWithNetwork(app, fixture.store, fakeExitProbe{}, func(hub.NetworkExitAccess) (string, error) { return "", nil })
response := do(app, http.MethodGet, "/api/browsers", "")
if response.Code != http.StatusOK {
t.Fatalf("reconcile returned %d: %s", response.Code, response.Body.String())
}
after, err := fixture.store.GetEnvironmentContext(context.Background(), "account-a")
if err != nil || after.RuntimeInstanceID == "" || after.RuntimeID != "container-id" {
t.Fatalf("reconcile did not activate replacement runtime: %#v err=%v", after, err)
}
var auditedRuntime string
var auditedBinding int64
if err := fixture.db.QueryRow(`
SELECT runtime_instance_id, binding_version FROM audit_event
WHERE event_type = 'environment_action_finished' AND action = 'reconcile'
ORDER BY id DESC LIMIT 1`).Scan(&auditedRuntime, &auditedBinding); err != nil {
t.Fatal(err)
}
if auditedRuntime != after.RuntimeInstanceID || auditedBinding != after.BindingVersion {
t.Fatalf("finished audit retained stale runtime: runtime=%q binding=%d current=%#v", auditedRuntime, auditedBinding, after)
}
}
func TestPostgresNonRunnableReconcileAuditsRuntimeRelease(t *testing.T) {
databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL")
if databaseURL == "" {
t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run PostgreSQL integration coverage")
}
ctx := context.Background()
for _, authorization := range []string{"authorized", "revoked"} {
for _, observation := range []string{"stopped", "missing"} {
for _, releaseFailure := range []bool{false, true} {
name := authorization + "/" + observation + "/success"
if releaseFailure {
name = authorization + "/" + observation + "/release failure"
}
t.Run(name, func(t *testing.T) {
fixture := newPostgresRebindFixture(t, databaseURL)
setFixtureAccountStatus(t, fixture.databaseURL, "active")
var err error
fixture.bound, err = fixture.store.ActivateRuntime(ctx, "account-a", "stopped-container",
fixture.bound.BindingVersion, fixture.bound.Exit.ID, "stopped-network")
if err != nil {
t.Fatal(err)
}
if _, err := fixture.db.ExecContext(ctx, `
UPDATE social_account SET status = 'paused', authorization_status = $1 WHERE id = 'account-a'`, authorization); err != nil {
t.Fatal(err)
}
if observation == "stopped" {
fixture.gateway.containers = []containerStatus{{
ID: "stopped-container", Alias: "account-a", State: "exited",
BindingVersion: fixture.bound.BindingVersion, NetworkExitID: fixture.bound.Exit.ID,
}}
}
var store hubStore = fixture.store
if releaseFailure {
store = failingRuntimeReleaseStore{Store: fixture.store, err: errors.New("release unavailable")}
}
app := fiber.New()
registerHubWithNetwork(app, store, fakeExitProbe{}, func(hub.NetworkExitAccess) (string, error) { return "", nil })
response := do(app, http.MethodGet, "/api/browsers", "")
wantStatus := http.StatusOK
wantOutcome, wantReason := "succeeded", "runtime_released"
if releaseFailure {
wantStatus, wantOutcome, wantReason = http.StatusInternalServerError, "failed", "runtime_release_failed"
}
if response.Code != wantStatus {
t.Fatalf("reconcile returned %d, want %d: %s", response.Code, wantStatus, response.Body.String())
}
after, err := fixture.store.GetEnvironmentContext(ctx, "account-a")
if err != nil || (after.RuntimeInstanceID != "") == !releaseFailure {
t.Fatalf("runtime lease after reconcile: %#v err=%v", after, err)
}
var events, operations, missingOperations int
var eventTypes, outcome, reason string
if err := fixture.db.QueryRowContext(ctx, `
SELECT count(*), count(DISTINCT operation_id),
count(*) FILTER (WHERE coalesce(operation_id, '') = ''),
string_agg(event_type, ',' ORDER BY id),
coalesce(max(outcome) FILTER (WHERE event_type = 'environment_action_finished'), ''),
coalesce(max(reason_code) FILTER (WHERE event_type = 'environment_action_finished'), '')
FROM audit_event WHERE action = 'reconcile'`).Scan(
&events, &operations, &missingOperations, &eventTypes, &outcome, &reason); err != nil {
t.Fatal(err)
}
if events != 2 || operations != 1 || missingOperations != 0 ||
eventTypes != "environment_action_requested,environment_action_finished" ||
outcome != wantOutcome || reason != wantReason {
t.Fatalf("reconcile audit mismatch: events=%d operations=%d missing=%d types=%q outcome=%q reason=%q",
events, operations, missingOperations, eventTypes, outcome, reason)
}
})
}
}
}
}
func TestPostgresImageDisableWaitsForEveryImageConsumer(t *testing.T) {
databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL")
if databaseURL == "" {
t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run PostgreSQL integration coverage")
}
for _, action := range []string{"create", "start", "reconcile", "rebind"} {
t.Run(action, func(t *testing.T) {
fixture := newPostgresRebindFixture(t, databaseURL)
if action == "start" || action == "reconcile" {
setFixtureAccountStatus(t, fixture.databaseURL, "active")
}
releaseLifecycle := make(chan struct{})
lifecycleStarted := make(chan struct{})
fixture.gateway.createStarted = lifecycleStarted
fixture.gateway.releaseCreate = releaseLifecycle
probe := networkExitProbe(fakeExitProbe{})
released := false
method, path, body, expected := http.MethodPost, "/api/browsers/account-a/"+action, "", http.StatusNoContent
switch action {
case "create":
if _, err := fixture.db.Exec(`DELETE FROM environment_binding; DELETE FROM browser_env`); err != nil {
t.Fatal(err)
}
path, expected = "/api/browsers", http.StatusCreated
body = fmt.Sprintf(`{"alias":"account-a","name":"甲","gateway":"gw-1","image_version":"148","fingerprint":{"seed":1},"account_id":"account-a","network_exit_id":"%s"}`, fixture.exit.ID)
case "reconcile":
method, path, expected = http.MethodGet, "/api/browsers", http.StatusOK
var err error
fixture.bound, err = fixture.store.ActivateRuntime(context.Background(), "account-a", "stale-container",
fixture.bound.BindingVersion, fixture.bound.Exit.ID, "network-old")
if err != nil {
t.Fatal(err)
}
fixture.gateway.containers = []containerStatus{{
ID: "stale-container", Alias: "account-a", State: "running", ProxyReady: true,
BindingVersion: fixture.bound.BindingVersion, NetworkExitID: "stale-exit", NetworkID: "network-old",
}}
case "rebind":
expected = http.StatusConflict
body = `{"network_exit_id":"` + fixture.exit.ID + `"}`
fixture.gateway.createStarted = nil
fixture.gateway.releaseCreate = nil
probe = &blockingExitProbe{started: lifecycleStarted, release: releaseLifecycle}
fixture.gateway.containers = []containerStatus{{
ID: "old-container", Alias: "account-a", State: "running", ProxyReady: true,
BindingVersion: fixture.bound.BindingVersion, NetworkExitID: fixture.bound.Exit.ID, NetworkID: "network-old",
}}
}
disableArrived := make(chan struct{})
var disableOnce sync.Once
app := fiber.New()
app.Use(func(c fiber.Ctx) error {
if c.Method() == http.MethodPut {
disableOnce.Do(func() { close(disableArrived) })
}
return c.Next()
})
registerHubWithNetwork(app, fixture.store, probe, func(hub.NetworkExitAccess) (string, error) { return "", nil })
server := httptest.NewServer(adaptor.FiberApp(app))
defer server.Close()
defer func() {
if !released {
close(releaseLifecycle)
}
}()
type result struct {
status int
err error
}
lifecycleDone := make(chan result, 1)
go func() {
request, err := http.NewRequest(method, server.URL+path, strings.NewReader(body))
if err == nil {
request.Header.Set("Content-Type", "application/json")
var response *http.Response
response, err = server.Client().Do(request)
if err == nil {
defer response.Body.Close()
lifecycleDone <- result{status: response.StatusCode}
return
}
}
lifecycleDone <- result{err: err}
}()
select {
case <-lifecycleStarted:
case result := <-lifecycleDone:
t.Fatalf("%s ended before lifecycle barrier: %#v", action, result)
case <-time.After(5 * time.Second):
t.Fatalf("%s did not reach lifecycle barrier", action)
}
disableDone := make(chan result, 1)
go func() {
request, err := http.NewRequest(http.MethodPut, server.URL+"/api/browser-images/148", strings.NewReader(
`{"image_ref":"registry.example/browser:148","enabled":false}`))
if err == nil {
request.Header.Set("Content-Type", "application/json")
var response *http.Response
response, err = server.Client().Do(request)
if err == nil {
defer response.Body.Close()
disableDone <- result{status: response.StatusCode}
return
}
}
disableDone <- result{err: err}
}()
<-disableArrived
select {
case result := <-disableDone:
t.Fatalf("image disable crossed in-flight %s: %#v", action, result)
case <-time.After(50 * time.Millisecond):
}
close(releaseLifecycle)
released = true
if result := <-lifecycleDone; result.err != nil || result.status != expected {
t.Fatalf("%s failed: %#v", action, result)
}
if result := <-disableDone; result.err != nil || result.status != http.StatusNoContent {
t.Fatalf("image disable after %s failed: %#v", action, result)
}
})
}
}
func TestImageDisableWaitsForUpgradeCommit(t *testing.T) {
releaseCreate := make(chan struct{})
gateway := &fakeGateway{
@@ -3650,6 +4003,29 @@ func TestStoppedReconcileReportsRuntimeReleaseFailure(t *testing.T) {
}
}
func TestListReconcileAuditsRuntimeReleaseFailure(t *testing.T) {
store := newMemoryStore()
store.envs["account-a"] = hub.Env{Alias: "account-a", Name: "甲", Gateway: "gw-1", ImageVersion: "148", Fingerprint: hub.Fingerprint{Seed: 1}}
store.bindings["account-a"] = hub.EnvironmentContext{
Env: store.envs["account-a"], AccountID: "account-a", BindingID: "account-a", BindingVersion: 1,
Exit: store.exits["exit-1"], RuntimeInstanceID: "runtime-instance", RuntimeID: "container-id",
}
store.releaseErr = errors.New("database unavailable")
gateway := &fakeGateway{token: "unit-test-gateway-token", containers: []containerStatus{{
ID: "container-id", Alias: "account-a", State: "exited", BindingVersion: 1, NetworkExitID: "exit-1",
}}}
app := newTestApp(t, store, gateway)
response := do(app, http.MethodGet, "/api/browsers", "")
if response.Code != http.StatusInternalServerError {
t.Fatalf("expected release failure, got %d: %s", response.Code, response.Body.String())
}
if len(store.actions) != 2 || store.actions[0].Action != "reconcile" || store.actions[1].Outcome != "failed" ||
store.actions[1].ReasonCode != "runtime_release_failed" {
t.Fatalf("reconcile release failure must have its own audit pair: %#v", store.actions)
}
}
func TestRebindRebuildsRunningContainerWithLatestBinding(t *testing.T) {
store := newMemoryStore()
store.exits["exit-2"] = hub.NetworkExit{ID: "exit-2", Protocol: "http", Host: "proxy.example", Port: 8080, HealthStatus: "healthy", Version: 1}
@@ -3790,6 +4166,69 @@ func TestDisableExitImmediatelyDiscardsRuntime(t *testing.T) {
}
}
func TestDisableExitWaitsForInFlightBrowserLifecycle(t *testing.T) {
releaseCreate := make(chan struct{})
gateway := &fakeGateway{
token: "unit-test-gateway-token",
createStarted: make(chan struct{}),
releaseCreate: releaseCreate,
}
store := newMemoryStore()
store.envs["account-a"] = hub.Env{Alias: "account-a", Name: "甲", Gateway: "gw-1", ImageVersion: "148", Fingerprint: hub.Fingerprint{Seed: 1}}
store.bindings["account-a"] = hub.EnvironmentContext{
Env: store.envs["account-a"], AccountID: "account-a", BindingID: "account-a", BindingVersion: 1, Exit: store.exits["exit-1"],
}
_ = store.CreateImage(nil, hub.Image{Version: "148", ImageRef: "registry.example/browser:148", Enabled: true})
gatewayServer := httptest.NewServer(gateway.handler(t))
defer gatewayServer.Close()
store.gateways["gw-1"] = hub.Gateway{Name: "gw-1", Endpoint: gatewayServer.URL, Token: gateway.token}
app := fiber.New()
registerHubWithNetwork(app, store, fakeExitProbe{}, func(hub.NetworkExitAccess) (string, error) { return "", nil })
server := httptest.NewServer(adaptor.FiberApp(app))
defer server.Close()
statuses := make(chan int, 2)
go func() {
response, err := server.Client().Post(server.URL+"/api/browsers/account-a/start", "application/json", nil)
if err != nil {
statuses <- 0
return
}
defer response.Body.Close()
statuses <- response.StatusCode
}()
select {
case <-gateway.createStarted:
case <-time.After(time.Second):
t.Fatal("start did not reach gateway create")
}
go func() {
response, err := server.Client().Post(server.URL+"/api/network-exits/exit-1/disable", "application/json", nil)
if err != nil {
statuses <- 0
return
}
defer response.Body.Close()
statuses <- response.StatusCode
}()
select {
case status := <-statuses:
t.Fatalf("lifecycle operation completed before the running create: %d", status)
case <-time.After(50 * time.Millisecond):
}
close(releaseCreate)
for range 2 {
if status := <-statuses; status != http.StatusNoContent && status != http.StatusOK {
t.Fatalf("unexpected lifecycle status: %d", status)
}
}
if store.exits["exit-1"].HealthStatus != "disabled" || store.bindings["account-a"].RuntimeID != "" {
t.Fatalf("disable did not clean the completed runtime: exit=%#v binding=%#v", store.exits["exit-1"], store.bindings["account-a"])
}
}
func TestDisableExitPropagatesUnknownGatewayReadAndPreservesLease(t *testing.T) {
for _, test := range []struct {
name string
+15 -6
View File
@@ -83,8 +83,11 @@ func registerPhaseA(app *fiber.App, store *phasea.Store, runtimeStore runtimeSto
})
app.Post("/api/phase-a/accounts/:id/pause", func(c fiber.Ctx) error {
runtimeOperations.Lock()
defer runtimeOperations.Unlock()
unlock, err := lockAccountResources(c.Context(), runtimeStore, c.Params("id"))
if err != nil {
return hubError(c, err)
}
defer unlock()
if err := store.PauseAccount(c.Context(), c.Params("id")); err != nil {
return phaseAError(c, err)
}
@@ -97,8 +100,11 @@ func registerPhaseA(app *fiber.App, store *phasea.Store, runtimeStore runtimeSto
})
app.Post("/api/phase-a/accounts/:id/resume", func(c fiber.Ctx) error {
runtimeOperations.Lock()
defer runtimeOperations.Unlock()
unlock, err := lockAccountResources(c.Context(), runtimeStore, c.Params("id"))
if err != nil {
return hubError(c, err)
}
defer unlock()
if err := store.ResumeAccount(c.Context(), c.Params("id")); err != nil {
if errors.Is(err, phasea.ErrConflict) {
return accountResumeConflict(c, store, runtimeStore, c.Params("id"))
@@ -109,8 +115,11 @@ func registerPhaseA(app *fiber.App, store *phasea.Store, runtimeStore runtimeSto
})
app.Post("/api/phase-a/accounts/:id/revoke", func(c fiber.Ctx) error {
runtimeOperations.Lock()
defer runtimeOperations.Unlock()
unlock, err := lockAccountResources(c.Context(), runtimeStore, c.Params("id"))
if err != nil {
return hubError(c, err)
}
defer unlock()
if err := store.RevokeAccount(c.Context(), c.Params("id")); err != nil {
return phaseAError(c, err)
}
+6 -2
View File
@@ -61,10 +61,14 @@ DOCKER_GID=$(stat -c %g /var/run/docker.sock) docker compose up --build
草稿经 `POST /api/phase-a/confirmations` 显式确认后才可投递到 `/api/phase-a/tasks`。任务由幂等键去重;`POST /api/phase-a/mock/execute` 使用 `FOR UPDATE SKIP LOCKED` 领取一分钟租约,执行前统一核对账号、草稿和确认版本。缺少确认或版本不一致会进入 `needs_confirmation`,暂停账号或 Mock 策略结果会进入 `policy_hold`,不确定结果与过期租约进入 `needs_confirmation`;这些状态都不会自动重试。`GET /api/phase-a/audit` 只导出账号、确认版本、尝试和结果等非秘密证据。
启动时控制面先应用 Phase A v1,再由 Hub runner 顺序应用 v2、v3、v4;每一步都在事务和 advisory lock 下前向执行。v3 保留旧表、列和历史记录,旧账号回填为 `platform=mock` 并暂停,仅账号 ID 与环境 alias 相同的记录自动建立 bindingv4 追加环境动作审计字段与索引。其余记录等待显式绑定。本阶段不提供破坏性自动回滚。
启动时控制面先应用 Phase A v1,再由 Hub runner 顺序应用 v2 至 v12;每一步都在事务和 advisory lock 下前向执行。v3 保留旧表、列和历史记录,旧账号回填为 `platform=mock` 并暂停,仅账号 ID 与环境 alias 相同的记录自动建立 binding;v4 追加环境动作审计字段与索引v5 清理持久 fingerprint 中的旧代理字段,v6 增加可重试的 runtime cleanup 状态,v7 为 runtime lease 增加 binding version 并回填可确定的既有记录,v8 至 v10 补齐 cleanup/runtime 的不可变 generation 与兼容约束,v11、v12 增加任务恢复状态并修复兼容约束。其余记录等待显式绑定。本阶段不提供破坏性自动回滚。
`POST /api/network-exits` 只接受协议、主机、端口、已有 `credential_reference: {id}` 和预期出口身份;新出口为 `unchecked`,由 `POST /api/network-exits/:id/check` 经实际代理链路变为 `healthy``unhealthy``disable` 不可被检查重新启用。credential reference 的 `reference_key` 不出现在 API、日志或审计中;OS Keyring/Secret Manager bridge 在控制面进程启动前注入 `CREATORHUB_CREDENTIAL_<SHA256(reference_key)>`(大写十六进制),值为请求期解析的 `username:password`,控制面不持久化解析值。
`POST /api/browsers` 必须同时给出 `account_id``network_exit_id`。环境创建、启动和升级都会重新检查出口身份,只有 `healthy` 才调用网关;控制面强制下发代理和 `disable_non_proxied_udp`fingerprint 中的代理字段会被拒绝。显式 `POST /api/browsers/:alias/rebind` 只允许 paused、无 executing task 且无活动 runtime 的账号。`DELETE /api/browsers/:alias` 回收容器但保留稳定 binding、环境和命名 Profile 卷,后续 create 复用它们。create/start/stop/upgrade/recycle 均写共享 operation ID 的 requested/finished 审计对;网关断连且无法调和时 outcome 为 `unknown`
解析后的出口凭据只存在于控制面单次请求和网关内存转发器中;Docker inspect、容器环境、标签、挂载、`Config.Cmd` 与进程参数只包含 `docker-gateway` 的无凭据本地代理地址。stopped 环境启动时先删除旧容器并确认 runtime lease 释放,再按当前 binding 重建;控制面每 20 秒及列表读取时调和网关,续租 running runtime、释放 stopped/missing runtime,过期 lease 也会在绑定事务中回收
解析后的出口凭据只存在于控制面单次请求和网关内存转发器中;Docker inspect、容器环境、标签、挂载、`Config.Cmd` 与进程参数只包含 `docker-gateway` 的无凭据本地代理地址。网关内存代理以 alias、binding version 和 exit ID 共同标识 generation;生命周期操作按 alias 串行,重启恢复或重建必须重新核对该 generation,旧出口代理不能被新容器复用
stopped 环境启动时先删除旧容器并确认 runtime lease 释放,再按当前 binding 重建;控制面每 20 秒及列表读取时调和网关,续租 running runtime、释放 stopped/missing runtime,过期 lease 也会在绑定事务中回收。控制面用 PostgreSQL advisory transaction lock 按 alias 协调多副本;每个 Store 最多允许 5 个锁会话占用 10 连接池的一半,为锁内数据库调用保留连接。create 同时锁定账号 ID、alias、请求出口和请求镜像;start、reconcile/rebuild、rebind 和 upgrade 锁定 alias、当前出口及当前镜像(upgrade 还锁目标镜像),拿锁后重新读取出口与镜像版本。账号 pause/resume/revoke 使用账号 ID 与当前 binding alias 加入同一协调域;镜像禁用、引用更新或账号状态变更不能穿透在途生命周期。
非法 upgrade/rebind 目标在进入 advisory lock key 前按公开格式校验;审计仅保留环境原有的非秘密资源关联,并以 `upgrade_input_rejected` / `rebind_input_rejected` 写同一 operation ID 的 requested/finished 对。reconcile 恢复或重建后会重新读取 context,finished 事件关联实际激活的 runtime instance、binding version 与出口;后台释放 runtime 的成功或失败也写独立的 `reconcile` 审计对。
+7 -7
View File
@@ -73,13 +73,13 @@ curl --fail --silent --show-error \
docker compose exec -T postgres \
psql -U creatorhub -d creatorhub -tAc \
'SELECT 1 FROM schema_migration WHERE version = 2;' \
| grep -qx 2
'SELECT 1 FROM schema_migration WHERE version = 12;' \
| grep -qx 1
docker compose ps
```
健康检查应成功,浏览器列表接口应返回 JSON,迁移查询当前应输出 `2`,三个 Compose 服务应为运行状态。然后访问 <http://127.0.0.1:8080>;修改过 `CREATORHUB_PORT` 时使用对应端口。
健康检查应成功,浏览器列表接口应返回 JSON,迁移查询当前应输出 `1`,三个 Compose 服务应为运行状态。然后访问 <http://127.0.0.1:8080>;修改过 `CREATORHUB_PORT` 时使用对应端口。
排障时读取结构化服务日志:
@@ -183,8 +183,8 @@ docker compose up --detach --build
< "$BACKUP_FILE"
docker compose exec -T postgres \
psql -U creatorhub -d creatorhub_restore_check -v ON_ERROR_STOP=1 -tAc \
'SELECT 1 FROM schema_migration WHERE version = 2;' \
| grep -qx 2
'SELECT 1 FROM schema_migration WHERE version = 12;' \
| grep -qx 1
docker compose exec -T postgres \
dropdb --force -U creatorhub creatorhub_restore_check
@@ -230,8 +230,8 @@ docker compose up --detach --build
< "$BACKUP_FILE"
docker compose exec -T postgres \
psql -U creatorhub -d creatorhub -v ON_ERROR_STOP=1 -tAc \
'SELECT 1 FROM schema_migration WHERE version = 2;' \
| grep -qx 2
'SELECT 1 FROM schema_migration WHERE version = 12;' \
| grep -qx 1
git switch --detach "$RESTORE_REV"
if ! docker compose up --detach --build; then
+61 -3
View File
@@ -11,6 +11,7 @@ import (
"fmt"
"net/url"
"regexp"
"sort"
"strings"
"time"
"unicode/utf8"
@@ -62,6 +63,8 @@ var (
func ValidImageVersion(version string) bool { return imageVersionPattern.MatchString(version) }
func ValidNetworkExitID(id string) bool { return exitIDPattern.MatchString(id) }
var (
aliasPattern = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,31}$`)
gatewayNamePattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$`)
@@ -71,8 +74,9 @@ var (
)
type Store struct {
db *sql.DB
notify taskstate.Notifier
db *sql.DB
lockAdmission chan struct{}
notify taskstate.Notifier
}
// Gateway 是平台注册的 docker-gateway 实例;Token 由平台生成,明文存储供页面复制(开发阶段约定)。
@@ -116,7 +120,9 @@ func Open(ctx context.Context, databaseURL string) (*Store, error) {
db.Close()
return nil, errors.New("connect to hub database")
}
store := &Store{db: db}
// LockResources keeps one connection until the lifecycle operation finishes.
// Admit at most half the pool so those operations can still open nested DB calls.
store := &Store{db: db, lockAdmission: make(chan struct{}, 5)}
if err := store.migrate(ctx); err != nil {
db.Close()
return nil, err
@@ -137,6 +143,58 @@ func (s *Store) notifyTransitions(transitions []taskstate.Transition) {
}
}
// LockResources serializes lifecycle state across control-plane replicas. The
// transaction carries no data changes; rolling it back only releases the locks.
func (s *Store) LockResources(ctx context.Context, aliases, exitIDs, imageVersions []string) (func(), error) {
for _, alias := range aliases {
if !aliasPattern.MatchString(alias) {
return nil, ErrInvalid
}
}
for _, id := range exitIDs {
if !exitIDPattern.MatchString(id) {
return nil, ErrInvalid
}
}
for _, version := range imageVersions {
if !imageVersionPattern.MatchString(version) {
return nil, ErrInvalid
}
}
if len(aliases)+len(exitIDs)+len(imageVersions) == 0 {
return func() {}, nil
}
select {
case s.lockAdmission <- struct{}{}:
case <-ctx.Done():
return nil, ctx.Err()
}
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
<-s.lockAdmission
return nil, errors.New("begin resource lock")
}
resources := []struct {
namespace int
keys []string
}{{1542738013, aliases}, {1542738015, exitIDs}, {1542738014, imageVersions}}
for _, resource := range resources {
keys := append([]string(nil), resource.keys...)
sort.Strings(keys)
for _, key := range keys {
if _, err := tx.ExecContext(ctx, `SELECT pg_advisory_xact_lock($1, hashtext($2))`, resource.namespace, key); err != nil {
_ = tx.Rollback()
<-s.lockAdmission
return nil, errors.New("lock resource")
}
}
}
return func() {
_ = tx.Rollback()
<-s.lockAdmission
}, nil
}
func (s *Store) migrate(ctx context.Context) error {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
+121
View File
@@ -4,16 +4,137 @@ import (
"context"
"encoding/json"
"errors"
"fmt"
"os"
"reflect"
"slices"
"strings"
"sync/atomic"
"testing"
"time"
"git.ipao.vip/rogee/creator-hub/internal/taskstate"
)
func TestEnvironmentLocksCoordinateAcrossStoreInstances(t *testing.T) {
databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL")
if databaseURL == "" {
t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run PostgreSQL integration coverage")
}
ctx := context.Background()
testURL := isolatedDatabaseURL(t, databaseURL)
first := openFullyMigratedHub(t, ctx, testURL)
t.Cleanup(func() { _ = first.Close() })
second, err := Open(ctx, testURL)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = second.Close() })
unlockFirst, err := first.LockResources(ctx, []string{"account-a"}, nil, nil)
if err != nil {
t.Fatal(err)
}
firstReleased := false
defer func() {
if !firstReleased {
unlockFirst()
}
}()
differentAlias, err := second.LockResources(ctx, []string{"account-b"}, nil, nil)
if err != nil {
t.Fatalf("different aliases must not share a lock: %v", err)
}
differentAlias()
acquired := make(chan func(), 1)
errors := make(chan error, 1)
started := make(chan struct{})
go func() {
close(started)
unlock, lockErr := second.LockResources(ctx, []string{"account-a"}, nil, nil)
if lockErr != nil {
errors <- lockErr
return
}
acquired <- unlock
}()
<-started
select {
case unlock := <-acquired:
unlock()
t.Fatal("same alias lock did not block across Store instances")
case err := <-errors:
t.Fatal(err)
case <-time.After(50 * time.Millisecond):
}
unlockFirst()
firstReleased = true
select {
case unlock := <-acquired:
unlock()
case err := <-errors:
t.Fatal(err)
case <-time.After(time.Second):
t.Fatal("same alias lock was not released")
}
}
func TestResourceLocksReserveConnectionsForLifecycleQueries(t *testing.T) {
databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL")
if databaseURL == "" {
t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run PostgreSQL integration coverage")
}
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
store := openFullyMigratedHub(t, ctx, isolatedDatabaseURL(t, databaseURL))
t.Cleanup(func() { _ = store.Close() })
const workers = 10
var acquired atomic.Int32
startQueries := make(chan struct{})
done := make(chan error, workers)
for worker := 0; worker < workers; worker++ {
go func(worker int) {
unlock, err := store.LockResources(ctx, []string{fmt.Sprintf("account-%d", worker)}, nil, nil)
if err != nil {
done <- err
return
}
defer unlock()
acquired.Add(1)
<-startQueries
_, err = store.ListEnvs(ctx)
done <- err
}(worker)
}
deadline := time.NewTimer(100 * time.Millisecond)
ticker := time.NewTicker(time.Millisecond)
for acquired.Load() < workers {
select {
case <-ticker.C:
case <-deadline.C:
goto release
}
}
release:
ticker.Stop()
if !deadline.Stop() {
select {
case <-deadline.C:
default:
}
}
close(startQueries)
for worker := 0; worker < workers; worker++ {
if err := <-done; err != nil {
t.Fatalf("locked lifecycle query %d did not complete: %v", worker, err)
}
}
}
func TestFingerprintArgsFollowUpstreamCommandLineContract(t *testing.T) {
full := Fingerprint{
Seed: 2024, Platform: "windows", PlatformVersion: "11.0.0",