From 024bf14448a19c4fca255465ec21af598dd681e1 Mon Sep 17 00:00:00 2001 From: Rogee Date: Mon, 31 Aug 2026 19:37:17 +0800 Subject: [PATCH] HH-807: harden runtime coordination (#21) --- cmd/control-plane/hub.go | 313 ++++++++++++++++-- cmd/control-plane/hub_test.go | 441 ++++++++++++++++++++++++- cmd/control-plane/phasea.go | 21 +- docs/architecture/container-control.md | 8 +- docs/deployment.md | 14 +- internal/hub/store.go | 64 +++- internal/hub/store_test.go | 121 +++++++ 7 files changed, 937 insertions(+), 45 deletions(-) diff --git a/cmd/control-plane/hub.go b/cmd/control-plane/hub.go index 33e32ea..8588039 100644 --- a/cmd/control-plane/hub.go +++ b/cmd/control-plane/hub.go @@ -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) diff --git a/cmd/control-plane/hub_test.go b/cmd/control-plane/hub_test.go index d768025..b436f95 100644 --- a/cmd/control-plane/hub_test.go +++ b/cmd/control-plane/hub_test.go @@ -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 diff --git a/cmd/control-plane/phasea.go b/cmd/control-plane/phasea.go index e9aa124..49593dd 100644 --- a/cmd/control-plane/phasea.go +++ b/cmd/control-plane/phasea.go @@ -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) } diff --git a/docs/architecture/container-control.md b/docs/architecture/container-control.md index 5a269e0..41cd052 100644 --- a/docs/architecture/container-control.md +++ b/docs/architecture/container-control.md @@ -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 相同的记录自动建立 binding;v4 只追加环境动作审计字段与索引。其余记录等待显式绑定。本阶段不提供破坏性自动回滚。 +启动时控制面先应用 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_`(大写十六进制),值为请求期解析的 `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` 审计对。 diff --git a/docs/deployment.md b/docs/deployment.md index 74049ca..92f513c 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -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 服务应为运行状态。然后访问 ;修改过 `CREATORHUB_PORT` 时使用对应端口。 +健康检查应成功,浏览器列表接口应返回 JSON,迁移查询当前应输出 `1`,三个 Compose 服务应为运行状态。然后访问 ;修改过 `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 diff --git a/internal/hub/store.go b/internal/hub/store.go index e0a5413..8296d68 100644 --- a/internal/hub/store.go +++ b/internal/hub/store.go @@ -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 { diff --git a/internal/hub/store_test.go b/internal/hub/store_test.go index bb83a79..c9b829a 100644 --- a/internal/hub/store_test.go +++ b/internal/hub/store_test.go @@ -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",