package api import ( "context" "errors" "reflect" "sync" "time" hub "git.ipao.vip/rogee/creator-hub/internal/environment" "github.com/sirupsen/logrus" ) var runtimeUseRenewInterval = 20 * time.Second type runtimeUseStore interface { AcquireRuntimeUse(context.Context, string, string, string, string) (hub.RuntimeUseLease, error) RenewRuntimeUse(context.Context, string) (hub.RuntimeUseLease, error) ReleaseRuntimeUse(context.Context, string) error } type runtimeUseHandle struct { store runtimeUseStore lease hub.RuntimeUseLease ctx context.Context cancel context.CancelFunc done chan struct{} mu sync.Mutex renewErr error closeOnce sync.Once closeErr error } func beginRuntimeUse(ctx context.Context, store any, alias, purpose, ownerID string) (context.Context, *runtimeUseHandle, error) { useStore, ok := store.(runtimeUseStore) if !ok || useStore == nil || (reflect.ValueOf(useStore).Kind() == reflect.Ptr && reflect.ValueOf(useStore).IsNil()) { return ctx, nil, errors.New("runtime use lease store is unavailable") } lease, err := useStore.AcquireRuntimeUse(ctx, alias, purpose, ownerID, "") if err != nil { return ctx, nil, err } useCtx, cancel := context.WithCancel(ctx) handle := &runtimeUseHandle{store: useStore, lease: lease, ctx: useCtx, cancel: cancel, done: make(chan struct{})} go handle.renew() return useCtx, handle, nil } func beginRuntimeUseForEnvironment(ctx context.Context, store any, environment hub.EnvironmentContext, purpose, ownerID string) (context.Context, *runtimeUseHandle, error) { if environment.Alias == "" || environment.RuntimeInstanceID == "" { return ctx, nil, hub.ErrConflict } useCtx, handle, err := beginRuntimeUse(ctx, store, environment.Alias, purpose, ownerID) if err != nil { return ctx, nil, err } if handle.lease.RuntimeInstanceID != environment.RuntimeInstanceID { return ctx, nil, errors.Join(hub.ErrConflict, handle.Close()) } return useCtx, handle, nil } func (h *runtimeUseHandle) renew() { ticker := time.NewTicker(runtimeUseRenewInterval) defer ticker.Stop() defer close(h.done) for { select { case <-h.ctx.Done(): return case <-ticker.C: if _, err := h.store.RenewRuntimeUse(h.ctx, h.lease.Token); err != nil { h.mu.Lock() h.renewErr = err h.mu.Unlock() logrus.WithError(err).WithFields(logrus.Fields{ "runtime_use_token": h.lease.Token, "runtime_use_purpose": h.lease.Purpose, }).Error("runtime use lease renewal failed") h.cancel() return } } } } func (h *runtimeUseHandle) Close() error { if h == nil { return nil } h.closeOnce.Do(func() { h.cancel() <-h.done releaseCtx, cancel := context.WithTimeout(context.WithoutCancel(h.ctx), 5*time.Second) defer cancel() releaseErr := h.store.ReleaseRuntimeUse(releaseCtx, h.lease.Token) h.mu.Lock() defer h.mu.Unlock() h.closeErr = errors.Join(h.renewErr, releaseErr) }) return h.closeErr }