From 64c21deb1e8f63b486eaa685ca9043b2e89e49fd Mon Sep 17 00:00:00 2001 From: xgxgx Date: Fri, 7 Aug 2026 21:01:15 +0800 Subject: [PATCH 01/10] =?UTF-8?q?feat(deployments):=20=E4=BD=BF=E7=94=A8?= =?UTF-8?q?=20River=20=E5=AE=9E=E7=8E=B0=E5=AE=9A=E6=97=B6=E8=B0=83?= =?UTF-8?q?=E5=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cmd/migrate/main.go | 4 + docs/design/be/deployments-api-contract.md | 43 +- go.mod | 25 +- go.sum | 44 +- internal/agents/handler.go | 33 +- internal/api/server.go | 5 +- internal/config/defaults.go | 9 + internal/db/agents.go | 34 +- internal/db/db.go | 21 + internal/db/deployment_mapper.go | 54 +- internal/db/deployment_mapper.xml | 119 +++- internal/db/deployment_mapper_test.go | 119 +++- internal/db/deployment_run_mapper.go | 4 +- internal/db/deployment_run_mapper.xml | 6 +- internal/db/deployments.go | 283 ++++++++-- .../00049_schedule_deployments_with_river.sql | 52 ++ internal/db/migrations_test.go | 22 + internal/db/webhooks.go | 67 ++- internal/deployments/cron.go | 192 +++++++ internal/deployments/cron_test.go | 91 +++ internal/deployments/execution.go | 75 +++ internal/deployments/handler.go | 527 ++++++++---------- internal/deployments/handler_contract_test.go | 13 + internal/deployments/resources.go | 2 +- internal/deployments/scheduler.go | 404 ++++++++++++++ internal/deployments/scheduler_test.go | 72 +++ internal/webhooks/enqueuer.go | 55 +- internal/webhooks/enqueuer_test.go | 24 + main.go | 24 + tests/deployments_api_test.go | 426 +++++++++++++- tests/uuid_boundary_postgres_test.go | 3 +- .../ManagedAgentsPage.resources.suite.tsx | 11 + .../ManagedAgentsPage.test-utils.tsx | 24 + .../managed-agents/resources/detail.tsx | 3 +- .../managed-agents/resources/model.tsx | 5 - web/src/features/managed-agents/types.ts | 3 +- 36 files changed, 2459 insertions(+), 439 deletions(-) create mode 100644 internal/db/migrations/00049_schedule_deployments_with_river.sql create mode 100644 internal/deployments/cron.go create mode 100644 internal/deployments/cron_test.go create mode 100644 internal/deployments/execution.go create mode 100644 internal/deployments/scheduler.go create mode 100644 internal/deployments/scheduler_test.go diff --git a/cmd/migrate/main.go b/cmd/migrate/main.go index 136009c0..d8cd4279 100644 --- a/cmd/migrate/main.go +++ b/cmd/migrate/main.go @@ -10,6 +10,7 @@ import ( "github.com/superduck-ai/open-managed-agents/internal/config" "github.com/superduck-ai/open-managed-agents/internal/db" + "github.com/superduck-ai/open-managed-agents/internal/deployments" "github.com/superduck-ai/open-managed-agents/internal/logging" ) @@ -46,6 +47,9 @@ func run(logger *slog.Logger) error { if err := database.Migrate(ctx); err != nil { return fmt.Errorf("migrate database: %w", err) } + if err := deployments.MigrateRiver(ctx, database, logger.With("component", "deployment_scheduler")); err != nil { + return fmt.Errorf("migrate River: %w", err) + } logger.Info("database migrations applied") return nil } diff --git a/docs/design/be/deployments-api-contract.md b/docs/design/be/deployments-api-contract.md index 6418e2b9..7341ad83 100644 --- a/docs/design/be/deployments-api-contract.md +++ b/docs/design/be/deployments-api-contract.md @@ -51,9 +51,50 @@ API 密钥请求必须携带 `anthropic-version: 2023-06-01`,并在 `anthropic 公开合同没有说明的模糊行为保持不变,包括 `limit=0`、`schedule:null`、默认列表顺序,以及未记录的错误状态和幂等行为。 -本次对齐仅涉及 HTTP/API 边界和手动运行的状态检查,不新增调度器、数据库迁移、Filestore 或 Sandbox 投影变更、自动暂停行为及其他运行时功能。 +## Scheduled Deployment 执行 + +OMA 使用 River `v0.42.0` 持久执行 schedule。River 官方 migrator 在应用 PostgreSQL database 的 `public` schema 中创建并升级 `river_job`、`river_queue`、`river_leader`、`river_notification` 和 `river_migration`;应用表仍由 Goose 管理。`river_migration` 持久记录已应用版本,进程启动只检查并应用缺失版本,不会重建 River 表。`cmd/migrate up` 和开发环境自动迁移会使用同一个数据库连接配置,依次推进两套 migration。River 内部表的 DDL 不复制到应用 migration,避免升级 River 时出现两套 schema 定义。 + +Cron 统一由 `internal/deployments` 的 schedule 组件解析: + +- 只接受五段 POSIX Cron 与有效 IANA timezone;拒绝 seconds/year、`L/W/#/?` 和 `@daily` shortcut。 +- `upcoming_runs_at` 返回最多五个不含 jitter 的名义 UTC 时刻,不再使用 366 天扫描上限,因此闰日计划有效。 +- spring-forward 不存在的墙上时刻不触发;fall-back 重复的墙上时刻触发两次。 +- 实际 River Job 只正向延后。jitter 窗口为相邻名义 occurrence 间隔的 15%,下限 5 秒、上限 9 分钟;窗口内 offset 由 Deployment ID 与名义时刻的稳定哈希决定。这是 OMA 内部选择,不是 Claude 公开的哈希算法。 + +每个 Deployment 持久化 `schedule_revision` 和 `next_scheduled_at`。create、明确修改或清除 schedule、pause、unpause、archive 都使旧 revision 的 Job 失效;PATCH 未携带 schedule 时不写这三个调度字段。schedule revision 在锁定 Deployment 的事务内由数据库原子递增,避免锁外旧快照覆盖 worker 已推进的游标。unpause 只从当前时间之后的下一个 occurrence 恢复,不补暂停期间的触发。worker 成功提交一个 Run 后推进到下一个名义 occurrence。create、明确修改 schedule 和 unpause 通过 Yourbatis 公开的 `SQLTx()` 将同一个事务交给 River `InsertTx`,使游标与 Job 一起提交或回滚。worker 推进游标和启动回填后的 Job 仍由每 30 秒一次的 reconciliation 补齐,入队使用 `ByArgs` 保持幂等。启动回填或 reconciliation 遇到单条确定性的存量 schedule 解析错误时记录并跳过该 Deployment;数据库或 River 基础设施错误仍使启动失败。 + +```mermaid +sequenceDiagram + participant API as Deployment API + participant AppDB as Application schema + participant River as public River tables + participant Worker as Scheduled worker + + API->>AppDB: 开启 Yourbatis 事务并保存调度游标 + API->>River: 使用同一事务幂等插入 Job + Note over AppDB,River: Deployment 与 Job 一起提交或回滚 + River->>Worker: 到达 jitter 后的 trigger_at + Worker->>AppDB: 锁定并校验 active/revision/next_scheduled_at + Worker->>AppDB: 原子写 Session 或失败 Run、推进/暂停游标并写 webhook outbox + Worker-->>River: occurrence 已完成 +``` + +`deployment_runs.trigger_type` 区分 manual 与 schedule,`scheduled_at` 保存 schedule Run 的名义时刻;部分唯一索引 `(deployment_uuid, scheduled_at) WHERE trigger_type = 'schedule'` 是 River at-least-once 下的最终幂等边界。API 的 `trigger_context` 由这两列生成,数据库不重复保存同义 JSON。Run 只表示 Session 创建成功或失败,不跟踪 Session 后续执行。 + +失败行为按 Claude 公开合同处理: + +- 根 Agent 归档会在同一数据库事务自动归档其 Deployment 并写入 `deployment.archived` outbox;根 Agent 在触发时已删除也会自动归档并原子写入该事件,且不生成 Run。 +- 其他引用或配置失败生成最终失败 Run。只有公开的 14 类 paused-reason error 会自动暂停;`session_rate_limited_error` 与 `session_creation_rejected_error` 不暂停,并继续下一个 occurrence。 +- 数据库或进程级失败发生在最终 Run 提交之前时交给 River 重试;已提交成功或失败 Run 后不重试当前 occurrence。 +- paused Deployment 仍允许 manual Run;manual Run 不发送 `deployment_run.*` webhook。 + +组织级最多保留 1,000 个未归档且 schedule 非空的 Deployment。创建以及从无 schedule 更新为有 schedule 时会先锁定 organization 并在事务内检查额度,避免并发越界。 + +Deployment 生命周期发送 `deployment.created/updated/paused/unpaused/archived`;scheduled Run 发送同一 Run ID 的 `deployment_run.started`,随后发送 `succeeded` 或 `failed`。自动暂停还发送以 Deployment ID 为资源 ID 的 `deployment.paused`。成功创建 Session 时继续发送现有 Session webhook。Scheduled Run 的这些 webhook delivery jobs 与 Run、Session、Deployment 状态和调度游标在同一个 Yourbatis 事务中提交;任一 outbox 写入失败都会回滚 occurrence 并交给 River 重试。 主要参考资料: - - +- diff --git a/go.mod b/go.mod index a74dc929..a32d2b2e 100644 --- a/go.mod +++ b/go.mod @@ -19,15 +19,18 @@ require ( github.com/modelcontextprotocol/go-sdk v1.6.1 github.com/pressly/goose/v3 v3.27.1 github.com/redis/go-redis/v9 v9.20.1 + github.com/riverqueue/river v0.42.0 + github.com/riverqueue/river/riverdriver/riverdatabasesql v0.42.0 + github.com/robfig/cron/v3 v3.0.1 github.com/samber/lo v1.53.0 github.com/standard-webhooks/standard-webhooks/libraries v0.0.1 github.com/superduck-ai/e2b-go-sdk v0.0.1 - github.com/superduck-ai/yourbatis v0.1.1 + github.com/superduck-ai/yourbatis v0.1.2 go.opentelemetry.io/proto/otlp v1.10.0 go.yaml.in/yaml/v3 v3.0.4 golang.org/x/net v0.56.0 - golang.org/x/sync v0.21.0 - golang.org/x/text v0.38.0 + golang.org/x/sync v0.22.0 + golang.org/x/text v0.40.0 google.golang.org/protobuf v1.36.11 ) @@ -43,6 +46,7 @@ require ( github.com/bahlo/generic-list-go v0.2.0 // indirect github.com/buger/jsonparser v1.1.2 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/davecgh/go-spew v1.1.1 // indirect github.com/google/jsonschema-go v0.4.3 // indirect github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 // indirect github.com/invopop/jsonschema v0.14.0 // indirect @@ -50,26 +54,33 @@ require ( github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect github.com/jackc/puddle/v2 v2.2.2 // indirect github.com/klauspost/cpuid/v2 v2.2.11 // indirect - github.com/kr/text v0.2.0 // indirect + github.com/lib/pq v1.12.3 // indirect github.com/mfridman/interpolate v0.0.2 // indirect github.com/pb33f/ordered-map/v2 v2.3.1 // indirect + github.com/pmezard/go-difflib v1.0.0 // indirect + github.com/riverqueue/river/riverdriver v0.42.0 // indirect + github.com/riverqueue/river/rivershared v0.42.0 // indirect + github.com/riverqueue/river/rivertype v0.42.0 // indirect github.com/rogpeppe/go-internal v1.15.0 // indirect github.com/segmentio/asm v1.2.1 // indirect github.com/segmentio/encoding v0.5.4 // indirect github.com/sethvargo/go-retry v0.3.0 // indirect - github.com/tidwall/gjson v1.18.0 // indirect - github.com/tidwall/match v1.1.1 // indirect + github.com/stretchr/testify v1.11.1 // indirect + github.com/tidwall/gjson v1.19.0 // indirect + github.com/tidwall/match v1.2.0 // indirect github.com/tidwall/pretty v1.2.1 // indirect github.com/tidwall/sjson v1.2.5 // indirect github.com/yosida95/uritemplate/v3 v3.0.2 // indirect go.uber.org/atomic v1.11.0 // indirect + go.uber.org/goleak v1.3.0 // indirect go.uber.org/multierr v1.11.0 // indirect go.yaml.in/yaml/v4 v4.0.0-rc.2 // indirect - golang.org/x/mod v0.37.0 // indirect + golang.org/x/mod v0.38.0 // indirect golang.org/x/oauth2 v0.35.0 // indirect golang.org/x/sys v0.46.0 // indirect golang.org/x/tools v0.47.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260209200024-4cfbd4190f57 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260420184626-e10c466a9529 // indirect google.golang.org/grpc v1.80.0 // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/go.sum b/go.sum index bbdad44a..6c256259 100644 --- a/go.sum +++ b/go.sum @@ -48,7 +48,6 @@ github.com/buger/jsonparser v1.1.2 h1:frqHqw7otoVbk5M8LlE/L7HTnIq2v9RX6EJ48i9AxJ github.com/buger/jsonparser v1.1.2/go.mod h1:6RYKKt7H4d4+iWqouImQ9R2FZql3VbhNgx27UK13J/0= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= -github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -78,6 +77,8 @@ github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= github.com/invopop/jsonschema v0.14.0 h1:MHQqLhvpNUZfw+hM3AZDYK7jxO8FZoQeQM77g8iyZjg= github.com/invopop/jsonschema v0.14.0/go.mod h1:ygm6C2EaVNMBDPpaPlnOA2pFAxBnxGjFlMZABxm9n2I= +github.com/jackc/pgerrcode v0.0.0-20240316143900-6e2875d9b438 h1:Dj0L5fhJ9F82ZJyVOmBx6msDp/kfd1t9GRfny/mfJA0= +github.com/jackc/pgerrcode v0.0.0-20240316143900-6e2875d9b438/go.mod h1:a/s9Lp5W7n/DD0VrVoyJ00FbP2ytTPDVOivvn2bMlds= github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= @@ -92,6 +93,8 @@ github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/lib/pq v1.12.3 h1:tTWxr2YLKwIvK90ZXEw8GP7UFHtcbTtty8zsI+YjrfQ= +github.com/lib/pq v1.12.3/go.mod h1:/p+8NSbOcwzAEI7wiMXFlgydTwcgTr3OSKMsD2BitpA= github.com/mattn/go-isatty v0.0.21 h1:xYae+lCNBP7QuW4PUnNG61ffM4hVIfm+zUzDuSzYLGs= github.com/mattn/go-isatty v0.0.21/go.mod h1:ZXfXG4SQHsB/w3ZeOYbR0PrPwLy+n6xiMrJlRFqopa4= github.com/mfridman/interpolate v0.0.2 h1:pnuTK7MQIxxFz1Gr+rjSIx9u7qVjf5VOoM/u6BbAxPY= @@ -110,6 +113,20 @@ github.com/redis/go-redis/v9 v9.20.1 h1:sfCU6A8P3dXbKyWes02uxA2baehGux9dZHfEKtsT github.com/redis/go-redis/v9 v9.20.1/go.mod h1:v/M13XI1PVCDcm01VtPFOADfZtHf8YW3baQf57KlIkA= github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE= github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo= +github.com/riverqueue/river v0.42.0 h1:WDx0F0loJg41XvkHVyQVr9PUN2p9SPLr2Tx7zs5ONlY= +github.com/riverqueue/river v0.42.0/go.mod h1:pD+hDP0ZW3SbuTwh0CXDlYQ/M3Q7HDntT3Q4lhfYugY= +github.com/riverqueue/river/riverdriver v0.42.0 h1:uf7CcKEtuHN+hPtuQfn5a+L4gvHJ7JlFFH/xNRQHfo4= +github.com/riverqueue/river/riverdriver v0.42.0/go.mod h1:b2IBlA29E3H233XwgbiJlezdoALSWhetZKt1LlBJEQU= +github.com/riverqueue/river/riverdriver/riverdatabasesql v0.42.0 h1:SH05h9K0jUz8nATb9mSxlu+tLw0X8wg+FLkej66BV2c= +github.com/riverqueue/river/riverdriver/riverdatabasesql v0.42.0/go.mod h1:YexDu481X/GKf7Gr4GNJ7K5M3GMUcxKh0oYxfHYJ0U4= +github.com/riverqueue/river/riverdriver/riverpgxv5 v0.42.0 h1:VQ165fd7oV0nBAQfEl52NyrBLWJo8oM7NY3u0z20mn4= +github.com/riverqueue/river/riverdriver/riverpgxv5 v0.42.0/go.mod h1:x+Yx1dcPLuriu8TqDHpAvZ5YQRJVnh9mxuxlqS0I6HA= +github.com/riverqueue/river/rivershared v0.42.0 h1:GZ+EXh3QbcSPUx8FlLcVpkv9MrOf8tMd0365K7AeR3g= +github.com/riverqueue/river/rivershared v0.42.0/go.mod h1:EThAIEr49dlUQFhVLJcQGKoMlnPKOq+UxdMi9jevVsk= +github.com/riverqueue/river/rivertype v0.42.0 h1:K7YEpq4WZFhMXi2b6M1mjLD4WwnJvI15gkVe6FREpTY= +github.com/riverqueue/river/rivertype v0.42.0/go.mod h1:D1Ad+EaZiaXbQbJcJcfeicXJMBKno0n6UcfKI5Q7DIQ= +github.com/robfig/cron/v3 v3.0.1 h1:WdRxkvbJztn8LMz/QEvLN5sBU+xKpSqwwUO1Pjr4qDs= +github.com/robfig/cron/v3 v3.0.1/go.mod h1:eQICP3HwyT7UooqI/z+Ov+PtYAWygg1TEWWzGIFLtro= github.com/rogpeppe/go-internal v1.15.0 h1:D0RCU5rMAp+SpgkiNdrjfJ+LX4J1M32V2NeCY7EJ6hc= github.com/rogpeppe/go-internal v1.15.0/go.mod h1:DrUVZyrJU+txYW5/1kwtXQSMFio52ZOxX7yM1VHvnxs= github.com/samber/lo v1.53.0 h1:t975lj2py4kJPQ6haz1QMgtId2gtmfktACxIXArw3HM= @@ -129,13 +146,14 @@ github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/superduck-ai/e2b-go-sdk v0.0.1 h1:yW23X/fEQvaPdJHhDnp8PgnsaIixn+UTVMt94X4Kr0s= github.com/superduck-ai/e2b-go-sdk v0.0.1/go.mod h1:dWlhv18vJamYp5gQ3ktmjjrX7tUaMe2OyTJXUo8EeXY= -github.com/superduck-ai/yourbatis v0.1.1 h1:iAEi8Hrx+p6MjS8ciI2kL8ueQ7FOrFDEphnN1NP3GSc= -github.com/superduck-ai/yourbatis v0.1.1/go.mod h1:BlCyyT1yfU2Zxya89rDf89keXqsdcwP6Q3PFscIjvig= +github.com/superduck-ai/yourbatis v0.1.2 h1:Fp/YSqfRXvOeBU0+UjJMnIxpMMfbyLKFa5lWyHkk3nU= +github.com/superduck-ai/yourbatis v0.1.2/go.mod h1:BlCyyT1yfU2Zxya89rDf89keXqsdcwP6Q3PFscIjvig= github.com/tidwall/gjson v1.14.2/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk= -github.com/tidwall/gjson v1.18.0 h1:FIDeeyB800efLX89e5a8Y0BNH+LOngJyGrIWxG2FKQY= -github.com/tidwall/gjson v1.18.0/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk= -github.com/tidwall/match v1.1.1 h1:+Ho715JplO36QYgwN9PGYNhgZvoUSc9X2c80KVTi+GA= +github.com/tidwall/gjson v1.19.0 h1:xwxm7n691Uf3u5OFjzngavjGTh55KX5q/9w9xHW88JU= +github.com/tidwall/gjson v1.19.0/go.mod h1:V37/opeE/JbLUOfH0QTXiNez2l0RUjYUhpT4szFQAfc= github.com/tidwall/match v1.1.1/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM= +github.com/tidwall/match v1.2.0 h1:0pt8FlkOwjN2fPt4bIl4BoNxb98gGHN2ObFEDkrfZnM= +github.com/tidwall/match v1.2.0/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM= github.com/tidwall/pretty v1.2.0/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= github.com/tidwall/pretty v1.2.1 h1:qjsOFOWWQl+N3RsoF5/ssm1pHmJJwhjlSbZ51I6wMl4= github.com/tidwall/pretty v1.2.1/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= @@ -161,24 +179,26 @@ go.opentelemetry.io/proto/otlp v1.10.0 h1:IQRWgT5srOCYfiWnpqUYz9CVmbO8bFmKcwYxpu go.opentelemetry.io/proto/otlp v1.10.0/go.mod h1:/CV4QoCR/S9yaPj8utp3lvQPoqMtxXdzn7ozvvozVqk= go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y= go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= go.yaml.in/yaml/v4 v4.0.0-rc.2 h1:/FrI8D64VSr4HtGIlUtlFMGsm7H7pWTbj6vOLVZcA6s= go.yaml.in/yaml/v4 v4.0.0-rc.2/go.mod h1:aZqd9kCMsGL7AuUv/m/PvWLdg5sjJsZ4oHDEnfPPfY0= -golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ= -golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0= +golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk= +golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40= golang.org/x/net v0.56.0 h1:Rw8j/hFzGvJUZwNBXnAtf5sVDVt+65SK2C7IxCxZt5o= golang.org/x/net v0.56.0/go.mod h1:D3Ku6r+V6JROoZK144D2XfMHFcMq/0zSfLelVTCFKec= golang.org/x/oauth2 v0.35.0 h1:Mv2mzuHuZuY2+bkyWXIHMfhNdJAdwW3FuWeCPYN5GVQ= golang.org/x/oauth2 v0.35.0/go.mod h1:lzm5WQJQwKZ3nwavOZ3IS5Aulzxi68dUSgRHujetwEA= -golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM= -golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek= +golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.46.0 h1:noSf2Fq6F8DBgS+LysIkx7rIExoNHJsxOAtPp4rthXw= golang.org/x/sys v0.46.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= -golang.org/x/text v0.38.0 h1:sXmwo9DwP3OK9EZ7PqAdaooSGozfl/3a6/xJcbzPRhE= -golang.org/x/text v0.38.0/go.mod h1:YXZt3QhHUKYT53r2lLKFIVi6Ao1jdzrTR/KQ09qyxF4= +golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs= +golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY= golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q= golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA= gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= diff --git a/internal/agents/handler.go b/internal/agents/handler.go index e72636ec..da75e629 100644 --- a/internal/agents/handler.go +++ b/internal/agents/handler.go @@ -19,6 +19,7 @@ import ( "github.com/superduck-ai/open-managed-agents/internal/ids" "github.com/superduck-ai/open-managed-agents/internal/logging" "github.com/superduck-ai/open-managed-agents/internal/modelmapping" + "github.com/superduck-ai/open-managed-agents/internal/webhooks" "github.com/go-chi/chi/v5" "github.com/google/uuid" @@ -31,10 +32,15 @@ const ( var customToolNamePattern = regexp.MustCompile(`^[A-Za-z0-9_-]{1,128}$`) type Handler struct { - cfg config.Config - db *db.DB - logger *slog.Logger - router chi.Router + cfg config.Config + db *db.DB + webhooks webhookEnqueuer + logger *slog.Logger + router chi.Router +} + +type webhookEnqueuer interface { + PrepareDeliveryEvent(webhooks.EnqueueInput, time.Time) (db.WebhookDeliveryEvent, error) } type agentResponse struct { @@ -85,9 +91,9 @@ type agentReference struct { Version int `json:"version"` } -func NewHandler(cfg config.Config, database *db.DB, logger *slog.Logger) *Handler { +func NewHandler(cfg config.Config, database *db.DB, webhookEvents webhookEnqueuer, logger *slog.Logger) *Handler { logger = logging.LoggerOrDefault(logger) - h := &Handler{cfg: cfg, db: database, logger: logger} + h := &Handler{cfg: cfg, db: database, webhooks: webhookEvents, logger: logger} router := chi.NewRouter() router.NotFound(notFound) router.MethodNotAllowed(notFound) @@ -382,7 +388,18 @@ func (h *Handler) archive(w http.ResponseWriter, r *http.Request, agentID string httpapi.WriteJSON(w, http.StatusOK, h.fixtureAgent(agentID, 1, true)) return } - record, err := h.db.ArchiveAgent(r.Context(), principal.WorkspaceUUID, agentID) + var buildEvent db.DeploymentArchiveEventBuilder + if h.webhooks != nil { + createdAt := time.Now().UTC() + buildEvent = func(deployment db.Deployment) (db.WebhookDeliveryEvent, error) { + return h.webhooks.PrepareDeliveryEvent(webhooks.EnqueueInput{ + WorkspaceUUID: principal.WorkspaceUUID, OrganizationUUID: principal.OrganizationUUID, + WorkspaceExternalID: principal.WorkspaceExternalID, EventType: "deployment.archived", + ResourceID: deployment.ExternalID, + }, createdAt) + } + } + archived, err := h.db.ArchiveAgent(r.Context(), principal.WorkspaceUUID, agentID, buildEvent) if err != nil { if errors.Is(err, db.ErrNotFound) { httpapi.WriteError(w, r, httpapi.NewError(http.StatusNotFound, "not_found_error", "Agent not found: "+agentID)) @@ -392,7 +409,7 @@ func (h *Handler) archive(w http.ResponseWriter, r *http.Request, agentID string httpapi.WriteError(w, r, httpapi.NewError(http.StatusInternalServerError, "api_error", "Could not archive agent")) return } - httpapi.WriteJSON(w, http.StatusOK, responseFromAgent(record)) + httpapi.WriteJSON(w, http.StatusOK, responseFromAgent(archived)) } func (h *Handler) versionsRoute(w http.ResponseWriter, r *http.Request) { diff --git a/internal/api/server.go b/internal/api/server.go index 6d19f9c1..de1cf1a0 100644 --- a/internal/api/server.go +++ b/internal/api/server.go @@ -78,6 +78,7 @@ type ServerDeps struct { SandboxTimeoutExtender codesessions.SandboxTimeoutExtender FilestoreCredentials *filestoreapi.TokenCredentials FilestoreService *filestoreapi.Service + DeploymentScheduler *deploymentsapi.DeploymentScheduler } // NewServer 用显式依赖组装 HTTP API Server。 @@ -109,10 +110,10 @@ func NewServer(deps ServerDeps) *Server { platformStore: platformStore, filestoreCredentials: deps.FilestoreCredentials, admin: adminapi.NewHandler(deps.Config, deps.DB, componentLogger("admin")), - agents: agents.NewHandler(deps.Config, deps.DB, componentLogger("agents")), + agents: agents.NewHandler(deps.Config, deps.DB, webhookEnqueuer, componentLogger("agents")), batch: batches.NewHandler(deps.Config, deps.DB, deps.ObjectStore, componentLogger("batches")), codeSessions: codesessions.NewHandler(deps.Config, codeSessionService, deps.SandboxTimeoutExtender, codeSessionLogger), - deployments: deploymentsapi.NewHandler(deps.DB, webhookEnqueuer, componentLogger("deployments")), + deployments: deploymentsapi.NewHandler(deps.DB, webhookEnqueuer, deps.DeploymentScheduler, componentLogger("deployments")), deploymentRuns: deploymentsapi.NewRunsHandler(deps.DB, componentLogger("deployment_runs")), envs: environments.NewHandler(deps.Config, deps.DB, componentLogger("environments")), files: files.NewHandler(deps.Config, deps.DB, deps.ObjectStore, componentLogger("files")), diff --git a/internal/config/defaults.go b/internal/config/defaults.go index 17f49e22..c24b33ac 100644 --- a/internal/config/defaults.go +++ b/internal/config/defaults.go @@ -96,6 +96,15 @@ func setDefaultSeedAPIKeys(cfg *Config) { func defaultWebhookEventTypes() []string { return []string{ + "deployment.created", + "deployment.updated", + "deployment.paused", + "deployment.unpaused", + "deployment.archived", + "deployment.deleted", + "deployment_run.started", + "deployment_run.succeeded", + "deployment_run.failed", "session.created", "session.pending", "session.running", diff --git a/internal/db/agents.go b/internal/db/agents.go index 102f9235..26729640 100644 --- a/internal/db/agents.go +++ b/internal/db/agents.go @@ -182,10 +182,36 @@ func isNullJSON(raw json.RawMessage) bool { return value == nil } -func (d *DB) ArchiveAgent(ctx context.Context, workspaceUUID string, externalID string) (Agent, error) { - mapper := NewAgentMapper(d.mapperDB) - row, err := mapper.ArchiveByExternalID(ctx, workspaceUUID, externalID) - return agentFromRow(row, err) +type DeploymentArchiveEventBuilder func(Deployment) (WebhookDeliveryEvent, error) + +func (d *DB) ArchiveAgent(ctx context.Context, workspaceUUID string, externalID string, buildEvent DeploymentArchiveEventBuilder) (Agent, error) { + var archived Agent + var deployments []Deployment + err := d.mapperDB.Transaction(ctx, func(executor yourbatis.Executor) error { + row, err := NewAgentMapper(executor).ArchiveByExternalID(ctx, workspaceUUID, externalID) + if err != nil { + return mapNoRows(err) + } + archived = row.agent() + rows, err := NewDeploymentMapper(executor).ArchiveByRootAgent(ctx, workspaceUUID, externalID) + if err != nil { + return err + } + deployments = deploymentsFromRows(rows) + if buildEvent == nil { + return nil + } + events := make([]WebhookDeliveryEvent, 0, len(deployments)) + for _, deployment := range deployments { + event, err := buildEvent(deployment) + if err != nil { + return err + } + events = append(events, event) + } + return enqueueWebhookDeliveryEventsTx(ctx, executor, workspaceUUID, events) + }) + return archived, err } func (d *DB) ListAgentsPage(ctx context.Context, params ListAgentsPageParams) ([]Agent, bool, error) { diff --git a/internal/db/db.go b/internal/db/db.go index e66fd3b6..e5794e7a 100644 --- a/internal/db/db.go +++ b/internal/db/db.go @@ -36,6 +36,7 @@ var ( ErrLimitExceeded = errors.New("limit exceeded") ErrFileInUse = errors.New("file is in use") ErrFileReferenceNotFound = errors.New("file reference not found") + ErrStaleSchedule = errors.New("stale deployment schedule") ) type DB struct { @@ -143,6 +144,26 @@ func (d *DB) Close() { } } +// Transaction executes fn in a Yourbatis transaction owned by DB. +func (d *DB) Transaction(ctx context.Context, fn func(*yourbatis.Tx) error) error { + if fn == nil { + return errors.New("database transaction callback is nil") + } + return d.mapperDB.Transaction(ctx, func(executor yourbatis.Executor) error { + tx, ok := executor.(*yourbatis.Tx) + if !ok { + return fmt.Errorf("unexpected Yourbatis transaction type %T", executor) + } + return fn(tx) + }) +} + +// RiverSQLDB exposes the shared database/sql wrapper only to the River queue +// integration. Application queries must continue to use Yourbatis mappers. +func (d *DB) RiverSQLDB() *sql.DB { + return d.sql +} + func EnsureDatabase(ctx context.Context, databaseURL string) error { var candidates []string for _, maintenanceDB := range []string{"postgres", "template1"} { diff --git a/internal/db/deployment_mapper.go b/internal/db/deployment_mapper.go index 9dafc09b..69d3a869 100644 --- a/internal/db/deployment_mapper.go +++ b/internal/db/deployment_mapper.go @@ -27,6 +27,8 @@ type deploymentMapperRow struct { ResourceSecrets []byte `db:"resource_secrets"` VaultIDs []byte `db:"vault_ids"` Schedule []byte `db:"schedule"` + ScheduleRevision int64 `db:"schedule_revision"` + NextScheduledAt *time.Time `db:"next_scheduled_at"` LastRunAt *time.Time `db:"last_run_at"` Status string `db:"status"` PausedReason []byte `db:"paused_reason"` @@ -56,6 +58,9 @@ type deploymentWriteParams struct { ResourceSecrets []byte VaultIDs []byte Schedule []byte + ScheduleChanged bool + ScheduleRevision int64 + NextScheduledAt *time.Time LastRunAt *time.Time Status string PausedReason []byte @@ -63,6 +68,45 @@ type deploymentWriteParams struct { UpdatedAt time.Time } +type unpauseDeploymentParams struct { + WorkspaceUUID string + ExternalID string + NextScheduledAt *time.Time +} + +type deploymentScheduleRow struct { + WorkspaceUUID string `db:"workspace_uuid"` + ExternalID string `db:"external_id"` + Schedule []byte `db:"schedule"` + ScheduleRevision int64 `db:"schedule_revision"` + NextScheduledAt *time.Time `db:"next_scheduled_at"` +} + +type setInitialNextScheduledAtParams struct { + WorkspaceUUID string + ExternalID string + ScheduleRevision int64 + NextScheduledAt time.Time +} + +type advanceDeploymentScheduleParams struct { + WorkspaceUUID string + ExternalID string + ScheduleRevision int64 + ScheduledAt time.Time + NextScheduledAt *time.Time + LastRunAt time.Time +} + +type pauseScheduledDeploymentParams struct { + WorkspaceUUID string + ExternalID string + ScheduleRevision int64 + ScheduledAt time.Time + PausedReason []byte + LastRunAt time.Time +} + type deploymentPageMapperParams struct { WorkspaceUUID string FetchLimit int @@ -76,12 +120,20 @@ type deploymentPageMapperParams struct { type DeploymentMapper interface { Insert(ctx context.Context, params deploymentWriteParams) (deploymentMapperRow, error) + LockOrganization(ctx context.Context, organizationUUID string) (string, error) + CountScheduledByOrganization(ctx context.Context, organizationUUID string) (int64, error) FindByExternalID(ctx context.Context, workspaceUUID, externalID string) (deploymentMapperRow, error) LockByExternalID(ctx context.Context, workspaceUUID, externalID string) (deploymentMapperRow, error) UpdateByExternalID(ctx context.Context, params deploymentWriteParams) (deploymentMapperRow, error) ArchiveByExternalID(ctx context.Context, workspaceUUID, externalID string) (deploymentMapperRow, error) + ArchiveByRootAgent(ctx context.Context, workspaceUUID, agentExternalID string) ([]deploymentMapperRow, error) PauseByExternalID(ctx context.Context, workspaceUUID, externalID string, pausedReason []byte) (deploymentMapperRow, error) - UnpauseByExternalID(ctx context.Context, workspaceUUID, externalID string) (deploymentMapperRow, error) + UnpauseByExternalID(ctx context.Context, params unpauseDeploymentParams) (deploymentMapperRow, error) + ListActiveSchedules(ctx context.Context) ([]deploymentScheduleRow, error) + ListSchedulesMissingNextScheduledAt(ctx context.Context) ([]deploymentScheduleRow, error) + SetInitialNextScheduledAt(ctx context.Context, params setInitialNextScheduledAtParams) (int64, error) + AdvanceSchedule(ctx context.Context, params advanceDeploymentScheduleParams) (int64, error) + PauseAfterScheduledRun(ctx context.Context, params pauseScheduledDeploymentParams) (int64, error) ListPage(ctx context.Context, params deploymentPageMapperParams) ([]deploymentMapperRow, error) UpdateLastRun(ctx context.Context, workspaceUUID, externalID string, lastRunAt time.Time) (int64, error) } diff --git a/internal/db/deployment_mapper.xml b/internal/db/deployment_mapper.xml index 6007721b..fce6cd22 100644 --- a/internal/db/deployment_mapper.xml +++ b/internal/db/deployment_mapper.xml @@ -20,6 +20,8 @@ resource_secrets, vault_ids, schedule, + schedule_revision, + next_scheduled_at, last_run_at, status, paused_reason, @@ -50,6 +52,8 @@ resource_secrets, vault_ids, schedule, + schedule_revision, + next_scheduled_at, last_run_at, status, paused_reason, @@ -75,6 +79,8 @@ CAST(#{params.ResourceSecrets,sensitive=true} AS jsonb), CAST(#{params.VaultIDs,sensitive=true} AS jsonb), CAST(#{params.Schedule,sensitive=true} AS jsonb), + #{params.ScheduleRevision}, + #{params.NextScheduledAt}, #{params.LastRunAt}, #{params.Status}, CAST(#{params.PausedReason,sensitive=true} AS jsonb), @@ -85,6 +91,22 @@ + + + + + SELECT workspace_uuid, external_id, schedule, schedule_revision, next_scheduled_at + FROM deployments + WHERE status = 'active' + AND archived_at IS NULL + AND deleted_at IS NULL + AND next_scheduled_at IS NOT NULL + ORDER BY next_scheduled_at, uuid + + + + + + UPDATE deployments + SET schedule_revision = schedule_revision + 1, + next_scheduled_at = #{params.NextScheduledAt}, + updated_at = NOW() + WHERE workspace_uuid = #{params.WorkspaceUUID} + AND external_id = #{params.ExternalID} + AND schedule_revision = #{params.ScheduleRevision} + AND status = 'active' + AND archived_at IS NULL + AND deleted_at IS NULL + AND schedule IS NOT NULL + AND next_scheduled_at IS NULL + + + + UPDATE deployments + SET last_run_at = #{params.LastRunAt}, + next_scheduled_at = #{params.NextScheduledAt}, + updated_at = #{params.LastRunAt} + WHERE workspace_uuid = #{params.WorkspaceUUID} + AND external_id = #{params.ExternalID} + AND schedule_revision = #{params.ScheduleRevision} + AND next_scheduled_at = #{params.ScheduledAt} + AND status = 'active' + AND archived_at IS NULL + AND deleted_at IS NULL + + + + UPDATE deployments + SET last_run_at = #{params.LastRunAt}, + status = 'paused', + paused_reason = CAST(#{params.PausedReason,sensitive=true} AS jsonb), + schedule_revision = schedule_revision + 1, + next_scheduled_at = NULL, + updated_at = #{params.LastRunAt} + WHERE workspace_uuid = #{params.WorkspaceUUID} + AND external_id = #{params.ExternalID} + AND schedule_revision = #{params.ScheduleRevision} + AND next_scheduled_at = #{params.ScheduledAt} + AND status = 'active' + AND archived_at IS NULL + AND deleted_at IS NULL + + + + - SELECT uuid - FROM organizations - WHERE uuid = #{organizationUUID} - FOR UPDATE - - - - - SELECT workspace_uuid, external_id, schedule, schedule_revision, next_scheduled_at - FROM deployments - WHERE status = 'active' - AND archived_at IS NULL - AND deleted_at IS NULL - AND next_scheduled_at IS NOT NULL - ORDER BY next_scheduled_at, uuid - - - + SELECT workspace_uuid, external_id, schedule, schedule_revision FROM deployments WHERE status = 'active' AND archived_at IS NULL AND deleted_at IS NULL AND schedule IS NOT NULL - AND next_scheduled_at IS NULL - ORDER BY uuid - - UPDATE deployments - SET schedule_revision = schedule_revision + 1, - next_scheduled_at = #{params.NextScheduledAt}, - updated_at = NOW() - WHERE workspace_uuid = #{params.WorkspaceUUID} - AND external_id = #{params.ExternalID} - AND schedule_revision = #{params.ScheduleRevision} - AND status = 'active' - AND archived_at IS NULL - AND deleted_at IS NULL - AND schedule IS NOT NULL - AND next_scheduled_at IS NULL - - - - UPDATE deployments - SET last_run_at = #{params.LastRunAt}, - next_scheduled_at = #{params.NextScheduledAt}, - updated_at = #{params.LastRunAt} - WHERE workspace_uuid = #{params.WorkspaceUUID} - AND external_id = #{params.ExternalID} - AND schedule_revision = #{params.ScheduleRevision} - AND next_scheduled_at = #{params.ScheduledAt} - AND status = 'active' - AND archived_at IS NULL - AND deleted_at IS NULL - - UPDATE deployments SET last_run_at = #{params.LastRunAt}, status = 'paused', paused_reason = CAST(#{params.PausedReason,sensitive=true} AS jsonb), schedule_revision = schedule_revision + 1, - next_scheduled_at = NULL, updated_at = #{params.LastRunAt} WHERE workspace_uuid = #{params.WorkspaceUUID} AND external_id = #{params.ExternalID} - AND schedule_revision = #{params.ScheduleRevision} - AND next_scheduled_at = #{params.ScheduledAt} - AND status = 'active' - AND archived_at IS NULL - AND deleted_at IS NULL - SELECT workspace_uuid, external_id, schedule, schedule_revision + SELECT workspace_uuid, external_id, schedule FROM deployments WHERE status = 'active' AND archived_at IS NULL @@ -208,7 +200,6 @@ SET last_run_at = #{params.LastRunAt}, status = 'paused', paused_reason = CAST(#{params.PausedReason,sensitive=true} AS jsonb), - schedule_revision = schedule_revision + 1, updated_at = #{params.LastRunAt} WHERE workspace_uuid = #{params.WorkspaceUUID} AND external_id = #{params.ExternalID} diff --git a/internal/db/deployment_mapper_test.go b/internal/db/deployment_mapper_test.go index 69f739af..89a82e13 100644 --- a/internal/db/deployment_mapper_test.go +++ b/internal/db/deployment_mapper_test.go @@ -15,9 +15,9 @@ import ( func TestDeploymentMapperBuilderContracts(t *testing.T) { now := time.Date(2026, time.August, 5, 1, 2, 3, 0, time.UTC) params := deploymentMapperTestWriteParams(now) - params.UpdateSchedule = true - withoutSchedule := params - withoutSchedule.UpdateSchedule = false + params.ScheduleChanged = true + unchangedSchedule := params + unchangedSchedule.ScheduleChanged = false page := deploymentPageMapperParams{ WorkspaceUUID: params.WorkspaceUUID, FetchLimit: 21, Cursor: &DeploymentPageCursor{CreatedAt: now, UUID: "00000000-0000-4000-8000-000000000009"}, @@ -37,8 +37,7 @@ func TestDeploymentMapperBuilderContracts(t *testing.T) { "params.CreatedByAPIKeyUUID", "params.EnvironmentUUID", "params.EnvironmentExternalID", "params.AgentUUID", "params.AgentExternalID", "params.AgentVersion", "params.AgentSnapshot", "params.Name", "params.Description", "params.Metadata", "params.InitialEvents", "params.Resources", - "params.ResourceSecrets", "params.VaultIDs", "params.Schedule", "params.ScheduleRevision", - "params.LastRunAt", "params.Status", + "params.ResourceSecrets", "params.VaultIDs", "params.Schedule", "params.LastRunAt", "params.Status", "params.PausedReason", "params.CreatedAt", "params.CreatedAt", }, wantSensitiveArgumentNames: deploymentSensitiveArgumentNames(true), @@ -74,11 +73,11 @@ func TestDeploymentMapperBuilderContracts(t *testing.T) { "params.UpdatedAt", "params.WorkspaceUUID", "params.ExternalID", }, wantSensitiveArgumentNames: deploymentSensitiveArgumentNames(false), - wantSQLFragments: []string{"UPDATE deployments", "schedule_revision = schedule_revision + 1", "workspace_uuid = $16", "RETURNING"}, + wantSQLFragments: []string{"UPDATE deployments", "schedule = CAST($14 AS jsonb)", "workspace_uuid = $16", "RETURNING"}, }}, - {"update without schedule", mapperBuilderContract{ + {"update without schedule change", mapperBuilderContract{ statement: deploymentMapperUpdateByExternalIDStatement, - bound: buildDeploymentMapperUpdateByExternalID(yourbatis.DialectPostgres, withoutSchedule), + bound: buildDeploymentMapperUpdateByExternalID(yourbatis.DialectPostgres, unchangedSchedule), wantID: "DeploymentMapper.UpdateByExternalID", wantKind: yourbatis.StatementUpdate, wantArgumentNames: []string{ "params.EnvironmentUUID", "params.EnvironmentExternalID", "params.AgentUUID", "params.AgentExternalID", @@ -103,7 +102,7 @@ func TestDeploymentMapperBuilderContracts(t *testing.T) { bound: buildDeploymentMapperArchiveByRootAgent(yourbatis.DialectPostgres, params.WorkspaceUUID, params.AgentExternalID), wantID: "DeploymentMapper.ArchiveByRootAgent", wantKind: yourbatis.StatementUpdate, wantArgumentNames: []string{"workspaceUUID", "agentExternalID"}, - wantSQLFragments: []string{"agent_external_id = $2", "schedule_revision = schedule_revision + 1"}, + wantSQLFragments: []string{"agent_external_id = $2", "archived_at = COALESCE"}, }}, {"pause", mapperBuilderContract{ statement: deploymentMapperPauseByExternalIDStatement, @@ -132,7 +131,7 @@ func TestDeploymentMapperBuilderContracts(t *testing.T) { bound: buildDeploymentMapperListActiveSchedules(yourbatis.DialectPostgres), wantID: "DeploymentMapper.ListActiveSchedules", wantKind: yourbatis.StatementSelect, wantArgumentNames: []string{}, - wantSQLFragments: []string{"SELECT workspace_uuid, external_id, schedule, schedule_revision", "schedule IS NOT NULL"}, + wantSQLFragments: []string{"SELECT workspace_uuid, external_id, schedule", "schedule IS NOT NULL"}, }}, {"pause after scheduled run", mapperBuilderContract{ statement: deploymentMapperPauseAfterScheduledRunStatement, @@ -160,9 +159,9 @@ func TestDeploymentMapperBuilderContracts(t *testing.T) { t.Run(test.name, func(t *testing.T) { assertMapperBuilderContract(t, test.contract) }) } - t.Run("update without schedule preserves revision", func(t *testing.T) { - bound := buildDeploymentMapperUpdateByExternalID(yourbatis.DialectPostgres, withoutSchedule) - if containsSQL(bound.SQL, "schedule_revision = schedule_revision + 1") { + t.Run("update without schedule change preserves schedule", func(t *testing.T) { + bound := buildDeploymentMapperUpdateByExternalID(yourbatis.DialectPostgres, unchangedSchedule) + if containsSQL(bound.SQL, "schedule =") { t.Fatalf("SQL unexpectedly changes schedule state: %q", bound.SQL) } }) @@ -303,7 +302,7 @@ func deploymentMapperTestWriteParams(now time.Time) deploymentWriteParams { AgentUUID: "00000000-0000-4000-8000-000000000006", AgentExternalID: "agent_test", AgentVersion: 1, AgentSnapshot: []byte(`{}`), Name: "test", Metadata: []byte(`{}`), InitialEvents: []byte(`[]`), Resources: []byte(`[]`), ResourceSecrets: []byte(`[]`), VaultIDs: []byte(`[]`), Schedule: []byte(`{}`), - ScheduleRevision: 1, Status: "active", PausedReason: []byte(`null`), CreatedAt: now, UpdatedAt: now, + Status: "active", PausedReason: []byte(`null`), CreatedAt: now, UpdatedAt: now, } } @@ -335,7 +334,7 @@ func deploymentMapperTestColumns() []string { "uuid", "external_id", "organization_uuid", "workspace_uuid", "created_by_api_key_uuid", "environment_uuid", "environment_external_id", "agent_uuid", "agent_external_id", "agent_version", "agent_snapshot", "name", "description", "metadata", "initial_events", "resources", "resource_secrets", - "vault_ids", "schedule", "schedule_revision", "last_run_at", "status", "paused_reason", "created_at", "updated_at", "archived_at", "deleted_at", + "vault_ids", "schedule", "last_run_at", "status", "paused_reason", "created_at", "updated_at", "archived_at", "deleted_at", } } @@ -346,7 +345,7 @@ func deploymentMapperTestRow() []driver.Value { "00000000-0000-4000-8000-000000000003", "00000000-0000-4000-8000-000000000004", "00000000-0000-4000-8000-000000000005", "env_test", "00000000-0000-4000-8000-000000000006", "agent_test", int64(1), []byte(`{}`), "test", nil, []byte(`{}`), []byte(`[]`), []byte(`[]`), []byte(`[]`), - []byte(`[]`), []byte(`{}`), int64(1), nil, "active", []byte(`null`), now, now, nil, nil, + []byte(`[]`), []byte(`{}`), nil, "active", []byte(`null`), now, now, nil, nil, } } diff --git a/internal/db/deployments.go b/internal/db/deployments.go index 7b8317a6..4345c811 100644 --- a/internal/db/deployments.go +++ b/internal/db/deployments.go @@ -34,7 +34,6 @@ type Deployment struct { ResourceSecrets json.RawMessage VaultIDs json.RawMessage Schedule json.RawMessage - ScheduleRevision int64 LastRunAt *time.Time Status string PausedReason json.RawMessage @@ -112,49 +111,39 @@ type UpdateDeploymentInput struct { } type DeploymentSchedule struct { - WorkspaceUUID string - ExternalID string - Schedule json.RawMessage - ScheduleRevision int64 + WorkspaceUUID string + ExternalID string + Schedule json.RawMessage } type ApplyScheduledOccurrenceInput struct { - WorkspaceUUID string - DeploymentExternalID string - ScheduleRevision int64 - ScheduledAt time.Time - Session *CreateSessionInput - Events []SessionEvent - Run DeploymentRun - AutoPauseReason json.RawMessage - ArchiveDeployment bool - Now time.Time + Deployment Deployment + ScheduledAt time.Time + Session *CreateSessionInput + Events []SessionEvent + Run DeploymentRun + AutoPauseReason json.RawMessage + ArchiveDeployment bool + Now time.Time } func (d *DB) CreateDeployment(ctx context.Context, deployment Deployment) (Deployment, error) { - var created Deployment - err := d.mapperDB.Transaction(ctx, func(executor yourbatis.Executor) error { - deploymentMapper := NewDeploymentMapper(executor) - deployment.ScheduleRevision = 0 - if len(deployment.Schedule) > 0 { - if err := checkScheduledDeploymentQuota(ctx, deploymentMapper, deployment.OrganizationUUID); err != nil { - return err - } - deployment.ScheduleRevision = 1 - } - row, err := deploymentMapper.Insert(ctx, deploymentWriteParamsFrom(deployment)) - if err != nil { - return err + mapper := NewDeploymentMapper(d.mapperDB) + if len(deployment.Schedule) > 0 { + if err := checkScheduledDeploymentQuota(ctx, mapper, deployment.OrganizationUUID); err != nil { + return Deployment{}, err } - created = row.deployment() - return nil - }) - return created, err + } + row, err := mapper.Insert(ctx, deploymentWriteParamsFrom(deployment)) + if err != nil { + return Deployment{}, err + } + return row.deployment(), nil } func (d *DB) GetDeployment(ctx context.Context, workspaceUUID string, externalID string) (Deployment, error) { - deploymentMapper := NewDeploymentMapper(d.mapperDB) - row, err := deploymentMapper.FindByExternalID(ctx, workspaceUUID, externalID) + mapper := NewDeploymentMapper(d.mapperDB) + row, err := mapper.FindByExternalID(ctx, workspaceUUID, externalID) if err != nil { return Deployment{}, mapNoRows(err) } @@ -164,8 +153,8 @@ func (d *DB) GetDeployment(ctx context.Context, workspaceUUID string, externalID func (d *DB) UpdateDeployment(ctx context.Context, workspaceUUID string, externalID string, input UpdateDeploymentInput) (Deployment, error) { var updated Deployment err := d.mapperDB.Transaction(ctx, func(executor yourbatis.Executor) error { - deploymentMapper := NewDeploymentMapper(executor) - current, err := deploymentMapper.LockByExternalID(ctx, workspaceUUID, externalID) + mapper := NewDeploymentMapper(executor) + current, err := mapper.LockByExternalID(ctx, workspaceUUID, externalID) if err != nil { return mapNoRows(err) } @@ -175,7 +164,7 @@ func (d *DB) UpdateDeployment(ctx context.Context, workspaceUUID string, externa next := input.Deployment scheduleChanged := input.ScheduleProvided && !sameJSON(current.Schedule, next.Schedule) if scheduleChanged && len(current.Schedule) == 0 && len(next.Schedule) > 0 { - if err := checkScheduledDeploymentQuota(ctx, deploymentMapper, current.OrganizationUUID); err != nil { + if err := checkScheduledDeploymentQuota(ctx, mapper, current.OrganizationUUID); err != nil { return err } } @@ -183,8 +172,8 @@ func (d *DB) UpdateDeployment(ctx context.Context, workspaceUUID string, externa params := deploymentWriteParamsFrom(next) params.WorkspaceUUID = workspaceUUID params.ExternalID = externalID - params.UpdateSchedule = scheduleChanged - row, err := deploymentMapper.UpdateByExternalID(ctx, params) + params.ScheduleChanged = scheduleChanged + row, err := mapper.UpdateByExternalID(ctx, params) if err != nil { return mapNoRows(err) } @@ -206,8 +195,8 @@ func checkScheduledDeploymentQuota(ctx context.Context, mapper DeploymentMapper, } func (d *DB) ArchiveDeployment(ctx context.Context, workspaceUUID string, externalID string) (Deployment, error) { - deploymentMapper := NewDeploymentMapper(d.mapperDB) - row, err := deploymentMapper.ArchiveByExternalID(ctx, workspaceUUID, externalID) + mapper := NewDeploymentMapper(d.mapperDB) + row, err := mapper.ArchiveByExternalID(ctx, workspaceUUID, externalID) if err != nil { return Deployment{}, mapNoRows(err) } @@ -215,15 +204,15 @@ func (d *DB) ArchiveDeployment(ctx context.Context, workspaceUUID string, extern } func (d *DB) PauseDeployment(ctx context.Context, workspaceUUID string, externalID string, pausedReason json.RawMessage) (Deployment, error) { - deploymentMapper := NewDeploymentMapper(d.mapperDB) - row, err := deploymentMapper.PauseByExternalID(ctx, workspaceUUID, externalID, agentJSONArg(pausedReason)) + mapper := NewDeploymentMapper(d.mapperDB) + row, err := mapper.PauseByExternalID(ctx, workspaceUUID, externalID, agentJSONArg(pausedReason)) if err != nil { return Deployment{}, mapNoRows(err) } return row.deployment(), nil } -func (d *DB) UnpauseDeployment(ctx context.Context, workspaceUUID, externalID string) (Deployment, error) { +func (d *DB) UnpauseDeployment(ctx context.Context, workspaceUUID string, externalID string) (Deployment, error) { mapper := NewDeploymentMapper(d.mapperDB) row, err := mapper.UnpauseByExternalID(ctx, workspaceUUID, externalID) if err != nil { @@ -236,8 +225,8 @@ func (d *DB) ListDeploymentsPage(ctx context.Context, params ListDeploymentsPage if params.Limit <= 0 { params.Limit = 20 } - deploymentMapper := NewDeploymentMapper(d.mapperDB) - rows, err := deploymentMapper.ListPage(ctx, deploymentPageParams(params)) + mapper := NewDeploymentMapper(d.mapperDB) + rows, err := mapper.ListPage(ctx, deploymentPageParams(params)) if err != nil { return nil, false, err } @@ -261,9 +250,7 @@ func (d *DB) CreateManualDeploymentRun(ctx context.Context, input CreateManualDe err := d.mapperDB.Transaction(ctx, func(executor yourbatis.Executor) error { deploymentMapper := NewDeploymentMapper(executor) runMapper := NewDeploymentRunMapper(executor) - deploymentRow, err := deploymentMapper.LockByExternalID( - ctx, input.Session.Session.WorkspaceUUID, input.DeploymentExternalID, - ) + deploymentRow, err := deploymentMapper.LockByExternalID(ctx, input.Session.Session.WorkspaceUUID, input.DeploymentExternalID) if err != nil { return mapNoRows(err) } @@ -297,80 +284,90 @@ func (d *DB) CreateManualDeploymentRun(ctx context.Context, input CreateManualDe return created, session, thread, events, err } -func (d *DB) ApplyScheduledOccurrenceTx(ctx context.Context, tx *yourbatis.Tx, input ApplyScheduledOccurrenceInput) error { - deploymentMapper := NewDeploymentMapper(tx) - runMapper := NewDeploymentRunMapper(tx) - row, err := deploymentMapper.LockByExternalID(ctx, input.WorkspaceUUID, input.DeploymentExternalID) - if err != nil { - return mapNoRows(err) - } - deployment := row.deployment() - if deployment.ArchivedAt != nil || deployment.Status != "active" || - deployment.ScheduleRevision != input.ScheduleRevision { - return ErrStaleSchedule - } - if input.ArchiveDeployment { - if _, err := deploymentMapper.ArchiveByExternalID(ctx, input.WorkspaceUUID, input.DeploymentExternalID); err != nil { - return err - } - return nil - } - if input.Session != nil { - workspace, err := NewAdminWorkspaceMapper(tx).FindByIdentifier( - ctx, deployment.OrganizationUUID, "", deployment.WorkspaceUUID, - ) +func (d *DB) ApplyScheduledOccurrence(ctx context.Context, input ApplyScheduledOccurrenceInput) error { + return d.mapperDB.Transaction(ctx, func(executor yourbatis.Executor) error { + deploymentMapper := NewDeploymentMapper(executor) + row, err := deploymentMapper.LockByExternalID(ctx, input.Deployment.WorkspaceUUID, input.Deployment.ExternalID) if err != nil { return mapNoRows(err) } - if workspace.ArchivedAt != nil { - return ErrWorkspaceArchived + deployment := row.deployment() + if deployment.ArchivedAt != nil || deployment.Status != "active" || + !sameJSON(deployment.Schedule, input.Deployment.Schedule) || + !sameDeploymentExecution(deployment, input.Deployment) { + return ErrStaleSchedule + } + if input.ArchiveDeployment { + if _, err := deploymentMapper.ArchiveByExternalID(ctx, deployment.WorkspaceUUID, deployment.ExternalID); err != nil { + return err + } + return nil + } + if input.Session != nil { + workspace, err := NewAdminWorkspaceMapper(executor).FindByIdentifier( + ctx, deployment.OrganizationUUID, "", deployment.WorkspaceUUID, + ) + if err != nil { + return mapNoRows(err) + } + if workspace.ArchivedAt != nil { + return ErrWorkspaceArchived + } } - } - run := deploymentRunFromDeployment(input.Run, deployment) - run.CreatedByAPIKeyUUID = deployment.CreatedByAPIKeyUUID - run.TriggerType = "schedule" - run.ScheduledAt = &input.ScheduledAt - run.CreatedAt = input.Now - if input.Session != nil { - session, _, _, _, err := insertSessionTx(ctx, tx, *input.Session) - if err != nil { - return err + runMapper := NewDeploymentRunMapper(executor) + run := deploymentRunFromDeployment(input.Run, deployment) + run.CreatedByAPIKeyUUID = deployment.CreatedByAPIKeyUUID + run.TriggerType = "schedule" + run.ScheduledAt = &input.ScheduledAt + run.CreatedAt = input.Now + if input.Session != nil { + session, _, _, _, err := insertSessionTx(ctx, executor, *input.Session) + if err != nil { + return err + } + if _, err = insertSessionEventsTx(ctx, executor, session, input.Events, false); err != nil { + return err + } + run.SessionExternalID = &session.ExternalID + run.Error = nil + } else { + run.SessionExternalID = nil } - if _, err = insertSessionEventsTx(ctx, tx, session, input.Events, false); err != nil { + _, err = runMapper.Insert(ctx, deploymentRunWriteParamsFrom(run)) + if err != nil { + if isUniqueViolationOnConstraint(err, "deployment_runs_schedule_occurrence_idx") { + return ErrStaleSchedule + } return err } - run.SessionExternalID = &session.ExternalID - run.Error = nil - } else { - run.SessionExternalID = nil - } - _, err = runMapper.Insert(ctx, deploymentRunWriteParamsFrom(run)) - if err != nil { - if isUniqueViolationOnConstraint(err, "deployment_runs_schedule_occurrence_idx") { - return ErrStaleSchedule + + if len(input.AutoPauseReason) > 0 { + _, err = deploymentMapper.PauseAfterScheduledRun(ctx, pauseScheduledDeploymentParams{ + WorkspaceUUID: deployment.WorkspaceUUID, ExternalID: deployment.ExternalID, + PausedReason: agentJSONArg(input.AutoPauseReason), LastRunAt: input.Now, + }) + } else { + _, err = deploymentMapper.UpdateLastRun( + ctx, deployment.WorkspaceUUID, deployment.ExternalID, input.Now, + ) } return err - } + }) +} - var rowsAffected int64 - if len(input.AutoPauseReason) > 0 { - rowsAffected, err = deploymentMapper.PauseAfterScheduledRun(ctx, pauseScheduledDeploymentParams{ - WorkspaceUUID: input.WorkspaceUUID, ExternalID: input.DeploymentExternalID, - PausedReason: agentJSONArg(input.AutoPauseReason), LastRunAt: input.Now, - }) - } else { - rowsAffected, err = deploymentMapper.UpdateLastRun( - ctx, input.WorkspaceUUID, input.DeploymentExternalID, input.Now, - ) - } - if err != nil { - return err - } - if rowsAffected != 1 { - return ErrStaleSchedule - } - return nil +func sameDeploymentExecution(left Deployment, right Deployment) bool { + return left.AgentUUID == right.AgentUUID && + left.AgentExternalID == right.AgentExternalID && + left.AgentVersion == right.AgentVersion && + sameJSON(left.AgentSnapshot, right.AgentSnapshot) && + left.EnvironmentUUID == right.EnvironmentUUID && + left.EnvironmentExternalID == right.EnvironmentExternalID && + sameJSON(left.Metadata, right.Metadata) && + sameJSON(left.InitialEvents, right.InitialEvents) && + sameJSON(left.Resources, right.Resources) && + sameJSON(left.ResourceSecrets, right.ResourceSecrets) && + sameJSON(left.VaultIDs, right.VaultIDs) } func (d *DB) CreateDeploymentRunFailure(ctx context.Context, deployment Deployment, run DeploymentRun) (DeploymentRun, error) { @@ -453,8 +450,7 @@ func deploymentWriteParamsFrom(deployment Deployment) deploymentWriteParams { Metadata: agentJSONArg(deployment.Metadata), InitialEvents: agentJSONArg(deployment.InitialEvents), Resources: agentJSONArg(deployment.Resources), ResourceSecrets: agentJSONArg(deployment.ResourceSecrets), VaultIDs: agentJSONArg(deployment.VaultIDs), Schedule: agentJSONArg(deployment.Schedule), - ScheduleRevision: deployment.ScheduleRevision, - LastRunAt: deployment.LastRunAt, Status: deployment.Status, PausedReason: agentJSONArg(deployment.PausedReason), + LastRunAt: deployment.LastRunAt, Status: deployment.Status, PausedReason: agentJSONArg(deployment.PausedReason), CreatedAt: deployment.CreatedAt, UpdatedAt: deployment.UpdatedAt, } } @@ -509,8 +505,7 @@ func (r deploymentMapperRow) deployment() Deployment { AgentSnapshot: bytes.Clone(r.AgentSnapshot), Name: r.Name, Description: r.Description, Metadata: bytes.Clone(r.Metadata), InitialEvents: bytes.Clone(r.InitialEvents), Resources: bytes.Clone(r.Resources), ResourceSecrets: bytes.Clone(r.ResourceSecrets), VaultIDs: bytes.Clone(r.VaultIDs), Schedule: bytes.Clone(r.Schedule), - ScheduleRevision: r.ScheduleRevision, - LastRunAt: r.LastRunAt, Status: r.Status, PausedReason: bytes.Clone(r.PausedReason), CreatedAt: r.CreatedAt, + LastRunAt: r.LastRunAt, Status: r.Status, PausedReason: bytes.Clone(r.PausedReason), CreatedAt: r.CreatedAt, UpdatedAt: r.UpdatedAt, ArchivedAt: r.ArchivedAt, DeletedAt: r.DeletedAt, } } diff --git a/internal/db/migrations/00050_schedule_deployments_with_river.sql b/internal/db/migrations/00050_schedule_deployments_with_river.sql index 1d8712ec..654151b6 100644 --- a/internal/db/migrations/00050_schedule_deployments_with_river.sql +++ b/internal/db/migrations/00050_schedule_deployments_with_river.sql @@ -1,7 +1,4 @@ -- +goose Up -alter table deployments - add column schedule_revision bigint not null default 0; - alter table deployment_runs add column scheduled_at timestamptz; @@ -40,6 +37,3 @@ end; alter table deployment_runs alter column trigger_context set not null, drop column scheduled_at; - -alter table deployments - drop column schedule_revision; diff --git a/internal/deployments/cron.go b/internal/deployments/cron.go index 7a758c87..467fa530 100644 --- a/internal/deployments/cron.go +++ b/internal/deployments/cron.go @@ -51,9 +51,9 @@ func parseDeploymentSchedule(raw json.RawMessage) (parsedSchedule, error) { return parsedSchedule{config: config, cron: cronSchedule}, nil } -func upcomingRuns(schedule cron.Schedule, now time.Time, archived bool) []string { +func upcomingRuns(schedule cron.Schedule, now time.Time, inactive bool) []string { values := make([]string, 0, upcomingRunCount) - if archived { + if inactive { return values } next := schedule.Next(now).UTC() diff --git a/internal/deployments/errors.go b/internal/deployments/errors.go index 583defaa..ab6826e7 100644 --- a/internal/deployments/errors.go +++ b/internal/deployments/errors.go @@ -84,3 +84,11 @@ func deploymentRunLoadError(err error, runID string) error { func deploymentFileMountConflict(cause error) error { return apperr.New(apperr.Conflict, "File resource mount_path conflicts with the session filesystem", cause) } + +func scheduledDeploymentLimitExceeded() error { + return apperr.New( + apperr.InvalidArgument, + fmt.Sprintf("an organization may have at most %d scheduled deployments", db.MaxScheduledDeploymentsPerOrganization), + db.ErrLimitExceeded, + ) +} diff --git a/internal/deployments/execution.go b/internal/deployments/execution.go index 4d19f919..ede4a9d4 100644 --- a/internal/deployments/execution.go +++ b/internal/deployments/execution.go @@ -17,7 +17,7 @@ func markRunPreparationRetryable(err error) error { return fmt.Errorf("%w: %v", errRetryableRunPreparation, err) } -type preparedDeploymentRun struct { +type preparedDeploymentExecution struct { RunID string Session db.CreateSessionInput Events []db.SessionEvent @@ -28,32 +28,36 @@ type deploymentSessionWork struct { Type string `json:"type"` } -func prepareDeploymentRun(deployment db.Deployment, now time.Time) (preparedDeploymentRun, error) { +func prepareDeploymentExecution( + deployment db.Deployment, + createdByAPIKeyUUID string, + now time.Time, +) (preparedDeploymentExecution, error) { sessionID, threadID, workID, runID, err := newRunIDs() if err != nil { - return preparedDeploymentRun{}, err + return preparedDeploymentExecution{}, err } events, outcomes, err := sessionEventsFromInitialEvents(deployment.InitialEvents, now) if err != nil { - return preparedDeploymentRun{}, err + return preparedDeploymentExecution{}, err } resources, err := sessionResourcesFromDeployment(deployment, now) if err != nil { - return preparedDeploymentRun{}, err + return preparedDeploymentExecution{}, err } deploymentID := deployment.ExternalID workData, err := jsonx.Encode(deploymentSessionWork{ID: sessionID, Type: "session"}) if err != nil { - return preparedDeploymentRun{}, err + return preparedDeploymentExecution{}, err } - return preparedDeploymentRun{ + return preparedDeploymentExecution{ RunID: runID, Events: events, Session: db.CreateSessionInput{ Session: db.Session{ UUID: uuid.NewString(), ExternalID: sessionID, OrganizationUUID: deployment.OrganizationUUID, WorkspaceUUID: deployment.WorkspaceUUID, - CreatedByAPIKeyUUID: deployment.CreatedByAPIKeyUUID, + CreatedByAPIKeyUUID: createdByAPIKeyUUID, EnvironmentUUID: deployment.EnvironmentUUID, EnvironmentExternalID: deployment.EnvironmentExternalID, AgentUUID: deployment.AgentUUID, AgentExternalID: deployment.AgentExternalID, AgentVersion: deployment.AgentVersion, AgentSnapshot: deployment.AgentSnapshot, diff --git a/internal/deployments/handler.go b/internal/deployments/handler.go index 2622c9e1..d3ce103f 100644 --- a/internal/deployments/handler.go +++ b/internal/deployments/handler.go @@ -234,10 +234,7 @@ type deploymentAgentSnapshot struct { func NewHandler(database *db.DB, webhookEvents webhookEnqueuer, logger *slog.Logger) *Handler { logger = logging.LoggerOrDefault(logger) - h := &Handler{ - db: database, webhooks: webhookEvents, - errorAdapter: httpapi.NewErrorAdapter(logger), - } + h := &Handler{db: database, webhooks: webhookEvents, errorAdapter: httpapi.NewErrorAdapter(logger)} wrap := h.errorAdapter.Wrap router := chi.NewRouter() router.NotFound(wrap(h.notFound)) @@ -362,7 +359,7 @@ func (h *Handler) create(w http.ResponseWriter, r *http.Request) error { return internalError("Could not generate deployment ID", fmt.Errorf("generate deployment ID: %w", err)) } now := time.Now().UTC() - deployment := db.Deployment{ + created, err := h.db.CreateDeployment(r.Context(), db.Deployment{ UUID: uuid.NewString(), ExternalID: deploymentID, OrganizationUUID: principal.OrganizationUUID, @@ -385,14 +382,10 @@ func (h *Handler) create(w http.ResponseWriter, r *http.Request) error { Status: "active", CreatedAt: now, UpdatedAt: now, - } - created, err := h.db.CreateDeployment(r.Context(), deployment) + }) if err != nil { if errors.Is(err, db.ErrLimitExceeded) { - return invalidRequest(fmt.Errorf( - "an organization may have at most %d scheduled deployments", - db.MaxScheduledDeploymentsPerOrganization, - )) + return scheduledDeploymentLimitExceeded() } return internalError("Could not create deployment", fmt.Errorf("create deployment %q: %w", deploymentID, err)) } @@ -557,38 +550,25 @@ func (h *Handler) updateRoute(w http.ResponseWriter, r *http.Request) error { return invalidRequest(err) } } - scheduleRaw := body.Schedule - if err := applyScheduleUpdate(&next, scheduleRaw); err != nil { - return invalidRequest(err) + scheduleProvided := body.Schedule != nil + if scheduleProvided { + next.Schedule, err = normalizeOptionalSchedule(body.Schedule) + if err != nil { + return invalidRequest(err) + } } next.UpdatedAt = time.Now().UTC() updated, err := h.db.UpdateDeployment(r.Context(), principal.WorkspaceUUID, deploymentID, db.UpdateDeploymentInput{ - Deployment: next, ScheduleProvided: len(scheduleRaw) > 0, + Deployment: next, ScheduleProvided: scheduleProvided, }) if err != nil { if errors.Is(err, db.ErrLimitExceeded) { - return invalidRequest(fmt.Errorf( - "an organization may have at most %d scheduled deployments", - db.MaxScheduledDeploymentsPerOrganization, - )) + return scheduledDeploymentLimitExceeded() } return deploymentLoadError(err, deploymentID) } return writeDeploymentResponse(w, updated) } - -func applyScheduleUpdate(next *db.Deployment, raw json.RawMessage) error { - if raw == nil { - return nil - } - schedule, err := normalizeOptionalSchedule(raw) - if err != nil { - return err - } - next.Schedule = schedule - return nil -} - func (h *Handler) archiveRoute(w http.ResponseWriter, r *http.Request) error { principal, err := requireAPIKey(r) if err != nil { @@ -644,13 +624,13 @@ func (h *Handler) runRoute(w http.ResponseWriter, r *http.Request) error { } referenceFailure, err := validateRunReferences(r.Context(), h.db, principal.WorkspaceUUID, deployment) if err != nil { - return deploymentLoadError(err, deploymentID) + referenceFailure = runError("unknown_error", "Could not create session") } if referenceFailure != nil { return h.writeRunReferenceFailure(w, r, principal, deployment, referenceFailure) } now := time.Now().UTC() - preparedRun, err := prepareDeploymentRun(deployment, now) + preparedRun, err := prepareDeploymentExecution(deployment, principal.APIKeyUUID, now) if err != nil { if errors.Is(err, errRetryableRunPreparation) { return deploymentLoadError(err, deploymentID) @@ -742,7 +722,7 @@ func validateRunReferences(ctx context.Context, database *db.DB, workspaceUUID s if agent.ArchivedAt != nil { return classifyReferenceFailure("agent", nil, true) } - return validateRunDependencies(ctx, database, workspaceUUID, deployment) + return validateSessionDependencies(ctx, database, workspaceUUID, deployment) } func validateRunDependencies(ctx context.Context, database *db.DB, workspaceUUID string, deployment db.Deployment) (*deploymentRunError, error) { @@ -778,6 +758,10 @@ func validateRunDependencies(ctx context.Context, database *db.DB, workspaceUUID return classifyReferenceFailure("skill", err, false) } } + return validateSessionDependencies(ctx, database, workspaceUUID, deployment) +} + +func validateSessionDependencies(ctx context.Context, database *db.DB, workspaceUUID string, deployment db.Deployment) (*deploymentRunError, error) { env, err := database.GetEnvironment(ctx, workspaceUUID, deployment.EnvironmentExternalID) if err != nil { return classifyReferenceFailure("environment", err, false) @@ -1124,18 +1108,23 @@ func sessionEventsFromInitialEvents(raw json.RawMessage, now time.Time) ([]db.Se return events, outcomesRaw, nil } -func scheduleResponse(scheduleRaw json.RawMessage, lastRunAt *time.Time, now time.Time, archived bool) (*deploymentScheduleResponse, error) { +func scheduleResponse(scheduleRaw json.RawMessage, lastRunAt *time.Time, now time.Time, inactive bool) *deploymentScheduleResponse { if len(scheduleRaw) == 0 || jsonx.IsNull(scheduleRaw) { - return nil, nil + return nil } - schedule, err := parseDeploymentSchedule(scheduleRaw) + config, err := jsonx.Decode[deploymentSchedule](scheduleRaw) if err != nil { - return nil, err + return nil } - return &deploymentScheduleResponse{ - deploymentSchedule: schedule.config, - LastRunAt: httpapi.OptionalTime(lastRunAt), UpcomingRunsAt: upcomingRuns(schedule.cron, now, archived), - }, nil + response := &deploymentScheduleResponse{ + deploymentSchedule: config, + LastRunAt: httpapi.OptionalTime(lastRunAt), + UpcomingRunsAt: []string{}, + } + if schedule, err := parseDeploymentSchedule(scheduleRaw); err == nil { + response.UpcomingRunsAt = upcomingRuns(schedule.cron, now, inactive) + } + return response } func writeDeploymentResponse(w http.ResponseWriter, deployment db.Deployment) error { @@ -1181,10 +1170,12 @@ func responseFromDeployment(deployment db.Deployment, now time.Time) (deployment if err != nil { return deploymentResponse{}, err } - schedule, err := scheduleResponse(deployment.Schedule, deployment.LastRunAt, now, deployment.ArchivedAt != nil) - if err != nil { - return deploymentResponse{}, err - } + schedule := scheduleResponse( + deployment.Schedule, + deployment.LastRunAt, + now, + deployment.Status != "active" || deployment.ArchivedAt != nil, + ) return deploymentResponse{ ID: deployment.ExternalID, Agent: agentReference(deployment.AgentExternalID, deployment.AgentVersion), @@ -1566,19 +1557,19 @@ func normalizeOutcomeRubric(raw json.RawMessage) (*deploymentOutcomeRubric, erro } func validateCheckout(raw json.RawMessage) error { - var request deploymentCheckoutRequest - if err := json.Unmarshal(raw, &request); err != nil { + var checkout deploymentCheckoutRequest + if err := json.Unmarshal(raw, &checkout); err != nil { return errors.New("checkout must be an object") } - checkoutType, err := parseRequiredRawString(request.Type, "type") + checkoutType, err := parseRequiredRawString(checkout.Type, "type") if err != nil { return err } switch checkoutType { case "branch": - _, err = parseRequiredRawString(request.Name, "name") + _, err = parseRequiredRawString(checkout.Name, "name") case "commit": - _, err = parseRequiredRawString(request.SHA, "sha") + _, err = parseRequiredRawString(checkout.SHA, "sha") default: err = errors.New("checkout.type must be branch or commit") } diff --git a/internal/deployments/handler_contract_test.go b/internal/deployments/handler_contract_test.go index 62a28b87..0164106b 100644 --- a/internal/deployments/handler_contract_test.go +++ b/internal/deployments/handler_contract_test.go @@ -49,6 +49,18 @@ func TestDeploymentRunResponseBuildsOfficialTriggerContext(t *testing.T) { } } +func TestScheduleResponseKeepsInvalidStoredCron(t *testing.T) { + response := scheduleResponse( + json.RawMessage(`{"type":"cron","expression":"bad","timezone":"UTC"}`), + nil, + time.Now(), + false, + ) + if response == nil || response.Expression != "bad" || len(response.UpcomingRunsAt) != 0 { + t.Fatalf("scheduleResponse() = %+v", response) + } +} + func TestDeploymentResourcesResponse(t *testing.T) { t.Run("rejects invalid stored resources", func(t *testing.T) { if _, err := deploymentResourcesResponse(json.RawMessage(`{"type":"file"}`)); err == nil { diff --git a/internal/deployments/scheduler.go b/internal/deployments/scheduler.go index d5f37f98..bc39d19b 100644 --- a/internal/deployments/scheduler.go +++ b/internal/deployments/scheduler.go @@ -17,7 +17,6 @@ import ( "github.com/superduck-ai/open-managed-agents/internal/db" "github.com/superduck-ai/open-managed-agents/internal/ids" "github.com/superduck-ai/open-managed-agents/internal/logging" - "github.com/superduck-ai/yourbatis" ) const ( @@ -26,12 +25,10 @@ const ( deploymentScheduleSyncInterval = 10 * time.Second ) -var errInvalidDeploymentSchedule = errors.New("invalid deployment schedule") - type scheduledDeploymentArgs struct { - WorkspaceUUID string `json:"workspace_uuid"` - DeploymentExternalID string `json:"deployment_id"` - ScheduleRevision int64 `json:"schedule_revision"` + WorkspaceUUID string `json:"workspace_uuid"` + DeploymentExternalID string `json:"deployment_id"` + Schedule deploymentSchedule `json:"schedule"` } func (scheduledDeploymentArgs) Kind() string { return "scheduled_deployment" } @@ -40,7 +37,7 @@ type DeploymentScheduler struct { database *db.DB client *river.Client[*sql.Tx] logger *slog.Logger - registered map[string]int64 + registered map[string]deploymentSchedule cancel context.CancelFunc done chan struct{} } @@ -76,7 +73,7 @@ func NewDeploymentScheduler(database *db.DB, logger *slog.Logger) (*DeploymentSc return nil, err } return &DeploymentScheduler{ - database: database, client: client, logger: logger, registered: make(map[string]int64), + database: database, client: client, logger: logger, registered: make(map[string]deploymentSchedule), }, nil } @@ -115,16 +112,29 @@ func (s *DeploymentScheduler) sync(ctx context.Context) error { } desired := make(map[string]struct{}, len(states)) for _, state := range states { - if err := s.update(state); err != nil { - if errors.Is(err, errInvalidDeploymentSchedule) { - state.Schedule = nil - _ = s.update(state) - s.logger.ErrorContext(ctx, "skip invalid stored deployment schedule", "deployment_id", state.ExternalID, "error", err) - continue - } - return fmt.Errorf("deployment %s: %w", state.ExternalID, err) + schedule, err := parseDeploymentSchedule(state.Schedule) + if err != nil { + s.client.PeriodicJobs().RemoveByID(state.ExternalID) + delete(s.registered, state.ExternalID) + s.logger.WarnContext(ctx, "skip invalid stored deployment schedule", "deployment_id", state.ExternalID, "error", err) + continue } desired[state.ExternalID] = struct{}{} + if registeredSchedule, ok := s.registered[state.ExternalID]; ok && registeredSchedule == schedule.config { + continue + } + job := river.NewPeriodicJob(schedule.cron, func() (river.JobArgs, *river.InsertOpts) { + return scheduledDeploymentArgs{ + WorkspaceUUID: state.WorkspaceUUID, DeploymentExternalID: state.ExternalID, + Schedule: schedule.config, + }, &river.InsertOpts{Queue: deploymentScheduleQueue} + }, &river.PeriodicJobOpts{ID: state.ExternalID}) + s.client.PeriodicJobs().RemoveByID(state.ExternalID) + if _, err := s.client.PeriodicJobs().AddSafely(job); err != nil { + delete(s.registered, state.ExternalID) + return fmt.Errorf("deployment %s: %w", state.ExternalID, err) + } + s.registered[state.ExternalID] = schedule.config } for id := range s.registered { if _, ok := desired[id]; !ok { @@ -135,34 +145,6 @@ func (s *DeploymentScheduler) sync(ctx context.Context) error { return nil } -func (s *DeploymentScheduler) update(state db.DeploymentSchedule) error { - if len(state.Schedule) == 0 { - s.client.PeriodicJobs().RemoveByID(state.ExternalID) - delete(s.registered, state.ExternalID) - return nil - } - if revision, ok := s.registered[state.ExternalID]; ok && revision == state.ScheduleRevision { - return nil - } - schedule, err := parseDeploymentSchedule(state.Schedule) - if err != nil { - return fmt.Errorf("%w: %v", errInvalidDeploymentSchedule, err) - } - job := river.NewPeriodicJob(schedule.cron, func() (river.JobArgs, *river.InsertOpts) { - return scheduledDeploymentArgs{ - WorkspaceUUID: state.WorkspaceUUID, DeploymentExternalID: state.ExternalID, - ScheduleRevision: state.ScheduleRevision, - }, &river.InsertOpts{Queue: deploymentScheduleQueue} - }, &river.PeriodicJobOpts{ID: state.ExternalID}) - s.client.PeriodicJobs().RemoveByID(state.ExternalID) - if _, err := s.client.PeriodicJobs().AddSafely(job); err != nil { - delete(s.registered, state.ExternalID) - return err - } - s.registered[state.ExternalID] = state.ScheduleRevision - return nil -} - func (s *DeploymentScheduler) syncLoop(ctx context.Context) { defer close(s.done) ticker := time.NewTicker(deploymentScheduleSyncInterval) @@ -194,25 +176,20 @@ func (w *scheduledDeploymentWorker) Work(ctx context.Context, job *river.Job[sch } return err } - if deployment.ArchivedAt != nil || deployment.Status != "active" || - deployment.ScheduleRevision != args.ScheduleRevision { + if deployment.ArchivedAt != nil || deployment.Status != "active" { + return nil + } + currentSchedule, err := parseDeploymentSchedule(deployment.Schedule) + if err != nil || currentSchedule.config != args.Schedule { return nil } now := time.Now().UTC() agent, agentErr := w.database.GetAgent(ctx, deployment.WorkspaceUUID, deployment.AgentExternalID) - if errors.Is(agentErr, db.ErrNotFound) || (agentErr == nil && agent.ArchivedAt != nil) { - err = w.applyOccurrence(ctx, job, db.ApplyScheduledOccurrenceInput{ - WorkspaceUUID: deployment.WorkspaceUUID, DeploymentExternalID: deployment.ExternalID, - ScheduleRevision: args.ScheduleRevision, ScheduledAt: scheduledAt, ArchiveDeployment: true, + if errors.Is(agentErr, db.ErrNotFound) || agent.ArchivedAt != nil { + return w.applyOccurrence(ctx, db.ApplyScheduledOccurrenceInput{ + Deployment: deployment, ScheduledAt: scheduledAt, ArchiveDeployment: true, }) - if errors.Is(err, db.ErrStaleSchedule) { - return nil - } - if err != nil { - return err - } - return nil } if agentErr != nil { return agentErr @@ -225,25 +202,21 @@ func (w *scheduledDeploymentWorker) Work(ctx context.Context, job *river.Job[sch if referenceFailure != nil { return w.recordFailure(ctx, job, deployment, referenceFailure, now) } - preparedRun, err := prepareDeploymentRun(deployment, now) + preparedRun, err := prepareDeploymentExecution(deployment, deployment.CreatedByAPIKeyUUID, now) if err != nil { if errors.Is(err, errRetryableRunPreparation) { return err } return w.recordFailure(ctx, job, deployment, runError("session_resource_not_found_error", err.Error()), now) } - err = w.applyOccurrence(ctx, job, db.ApplyScheduledOccurrenceInput{ - WorkspaceUUID: deployment.WorkspaceUUID, DeploymentExternalID: deployment.ExternalID, - ScheduleRevision: args.ScheduleRevision, ScheduledAt: scheduledAt, + err = w.applyOccurrence(ctx, db.ApplyScheduledOccurrenceInput{ + Deployment: deployment, ScheduledAt: scheduledAt, Session: &preparedRun.Session, Events: preparedRun.Events, Run: db.DeploymentRun{ UUID: uuid.NewString(), ExternalID: preparedRun.RunID, }, Now: now, }) - if errors.Is(err, db.ErrStaleSchedule) { - return nil - } if errors.Is(err, db.ErrWorkspaceArchived) { return w.recordFailure(ctx, job, deployment, runError("workspace_archived_error", "Workspace is archived"), now) } @@ -263,7 +236,6 @@ func (w *scheduledDeploymentWorker) recordFailure( failure *deploymentRunError, now time.Time, ) error { - args := job.Args scheduledAt := job.ScheduledAt.UTC() runID, err := ids.New("drun_") if err != nil { @@ -280,38 +252,27 @@ func (w *scheduledDeploymentWorker) recordFailure( return err } } - err = w.applyOccurrence(ctx, job, db.ApplyScheduledOccurrenceInput{ - WorkspaceUUID: deployment.WorkspaceUUID, DeploymentExternalID: deployment.ExternalID, - ScheduleRevision: args.ScheduleRevision, ScheduledAt: scheduledAt, + return w.applyOccurrence(ctx, db.ApplyScheduledOccurrenceInput{ + Deployment: deployment, ScheduledAt: scheduledAt, Run: db.DeploymentRun{ UUID: uuid.NewString(), ExternalID: runID, Error: runErrorJSON, }, AutoPauseReason: pausedReasonJSON, Now: now, }) - if errors.Is(err, db.ErrStaleSchedule) { - return nil - } - return err } func (w *scheduledDeploymentWorker) applyOccurrence( ctx context.Context, - job *river.Job[scheduledDeploymentArgs], input db.ApplyScheduledOccurrenceInput, ) error { - return w.database.Transaction(ctx, func(tx *yourbatis.Tx) error { - if err := w.database.ApplyScheduledOccurrenceTx(ctx, tx, input); err != nil { - return err - } - _, err := river.JobCompleteTx[*riverdatabasesql.Driver](ctx, tx.SQLTx(), job) - return err - }) + err := w.database.ApplyScheduledOccurrence(ctx, input) + if errors.Is(err, db.ErrStaleSchedule) { + return nil + } + return err } func shouldAutoPause(runError *deploymentRunError) bool { - if runError == nil { - return false - } switch runError.Type { case "environment_archived_error", "agent_archived_error", diff --git a/internal/deployments/scheduler_test.go b/internal/deployments/scheduler_test.go index fca28de0..0bb895b7 100644 --- a/internal/deployments/scheduler_test.go +++ b/internal/deployments/scheduler_test.go @@ -47,7 +47,4 @@ func TestShouldAutoPauseUsesOfficialAllowlist(t *testing.T) { t.Errorf("shouldAutoPause(%q) = true", errorType) } } - if shouldAutoPause(nil) { - t.Error("shouldAutoPause(nil) = true") - } } diff --git a/tests/deployments_api_test.go b/tests/deployments_api_test.go index f22b3a3b..09592d9c 100644 --- a/tests/deployments_api_test.go +++ b/tests/deployments_api_test.go @@ -15,7 +15,6 @@ import ( "github.com/google/uuid" "github.com/superduck-ai/open-managed-agents/internal/db" deploymentsapi "github.com/superduck-ai/open-managed-agents/internal/deployments" - "github.com/superduck-ai/yourbatis" ) type deploymentAPIResponse struct { @@ -366,54 +365,35 @@ func TestDeploymentsAPI(t *testing.T) { }) - t.Run("schedule revision changes only when schedule changes", func(t *testing.T) { - agent := createAgent(t, app, `{"model":"claude-opus-4-6","name":"deployments-execution-revision-agent"}`) + t.Run("failure changed schedule rejects an old occurrence", func(t *testing.T) { + agent := createAgent(t, app, `{"model":"claude-opus-4-6","name":"deployments-stale-schedule-agent"}`) defer cleanupAgentRows(t, app.pool, agent.ID) - env := createEnvironment(t, app, `{"name":"deployments-execution-revision-env"}`) + env := createEnvironment(t, app, `{"name":"deployments-stale-schedule-env"}`) defer cleanupEnvironmentRows(t, app.pool, env.ID) - created := createDeployment(t, app, `{ - "agent":`+quoteJSON(agent.ID)+`, - "environment_id":`+quoteJSON(env.ID)+`, - "name":"execution revision", - "initial_events":[{"type":"user.message","content":[{"type":"text","text":"before"}]}], - "schedule":{"type":"cron","expression":"*/10 * * * *","timezone":"UTC"} - }`) + created := createDeployment(t, app, deploymentBodyWithExtra(agent.ID, env.ID, `"schedule":{"type":"cron","expression":"*/10 * * * *","timezone":"UTC"}`)) defer cleanupDeploymentRows(t, app, created.ID) ctx := context.Background() ids := getDefaultDBIDs(t, app.pool) - original, err := app.db.GetDeployment(ctx, ids.WorkspaceUUID, created.ID) + deployment, err := app.db.GetDeployment(ctx, ids.WorkspaceUUID, created.ID) if err != nil { t.Fatalf("load scheduled deployment: %v", err) } - updateDeployment(t, app, created.ID, `{"initial_events":[{"type":"user.message","content":[{"type":"text","text":"after"}]}]}`) - current, err := app.db.GetDeployment(ctx, ids.WorkspaceUUID, created.ID) - if err != nil { - t.Fatalf("load updated deployment: %v", err) - } - if original.ScheduleRevision != 1 { - t.Fatalf("initial schedule_revision = %d, want 1", original.ScheduleRevision) - } - if current.ScheduleRevision != original.ScheduleRevision { - t.Fatalf("execution update schedule_revision = %d, want %d", current.ScheduleRevision, original.ScheduleRevision) - } - - updateDeployment(t, app, created.ID, `{"schedule":{"type":"cron","expression":"*/10 * * * *","timezone":"UTC"}}`) - unchanged, err := app.db.GetDeployment(ctx, ids.WorkspaceUUID, created.ID) - if err != nil { - t.Fatalf("load deployment after unchanged schedule: %v", err) - } - if unchanged.ScheduleRevision != original.ScheduleRevision { - t.Fatalf("unchanged schedule_revision = %d, want %d", unchanged.ScheduleRevision, original.ScheduleRevision) - } - updateDeployment(t, app, created.ID, `{"schedule":{"type":"cron","expression":"*/15 * * * *","timezone":"UTC"}}`) - changed, err := app.db.GetDeployment(ctx, ids.WorkspaceUUID, created.ID) - if err != nil { - t.Fatalf("load deployment after schedule change: %v", err) + scheduledAt := time.Now().UTC().Truncate(time.Minute) + err = applyScheduledOccurrence(ctx, app.db, db.ApplyScheduledOccurrenceInput{ + Deployment: deployment, ScheduledAt: scheduledAt, + Run: db.DeploymentRun{UUID: uuid.NewString(), ExternalID: "drun_stale_" + uuid.NewString()}, + Now: scheduledAt, + }) + if !errors.Is(err, db.ErrStaleSchedule) { + t.Fatalf("apply old scheduled occurrence error = %v, want ErrStaleSchedule", err) } - if changed.ScheduleRevision != original.ScheduleRevision+1 { - t.Fatalf("changed schedule_revision = %d, want %d", changed.ScheduleRevision, original.ScheduleRevision+1) + runs, _, err := app.db.ListDeploymentRunsPage(ctx, db.ListDeploymentRunsPageParams{ + WorkspaceUUID: ids.WorkspaceUUID, DeploymentExternalID: created.ID, Limit: 10, + }) + if err != nil || len(runs) != 0 { + t.Fatalf("runs after old occurrence = (%d, %v), want none", len(runs), err) } }) @@ -433,8 +413,7 @@ func TestDeploymentsAPI(t *testing.T) { } scheduledAt := time.Now().UTC().Truncate(time.Minute) input := db.ApplyScheduledOccurrenceInput{ - WorkspaceUUID: ids.WorkspaceUUID, DeploymentExternalID: created.ID, - ScheduleRevision: deployment.ScheduleRevision, ScheduledAt: scheduledAt, + Deployment: deployment, ScheduledAt: scheduledAt, Run: db.DeploymentRun{ UUID: uuid.NewString(), ExternalID: "drun_periodic_" + uuid.NewString(), Error: json.RawMessage(`{"type":"unknown_error","message":"test"}`), @@ -486,15 +465,14 @@ func TestDeploymentsAPI(t *testing.T) { }() err = applyScheduledOccurrence(ctx, app.db, db.ApplyScheduledOccurrenceInput{ - WorkspaceUUID: ids.WorkspaceUUID, DeploymentExternalID: created.ID, - ScheduleRevision: deployment.ScheduleRevision, ScheduledAt: scheduledAt, + Deployment: deployment, ScheduledAt: scheduledAt, Session: &db.CreateSessionInput{}, }) if !errors.Is(err, db.ErrWorkspaceArchived) { t.Fatalf("ApplyScheduledOccurrence() error = %v, want ErrWorkspaceArchived", err) } after, loadErr := app.db.GetDeployment(ctx, ids.WorkspaceUUID, created.ID) - if loadErr != nil || after.Status != "active" || after.ScheduleRevision != deployment.ScheduleRevision || after.LastRunAt != nil { + if loadErr != nil || after.Status != "active" || after.LastRunAt != nil { t.Fatalf("deployment after rejected occurrence = (%+v, %v), want unchanged", after, loadErr) } }) @@ -1065,9 +1043,7 @@ func cleanupDeploymentRows(t *testing.T, app *testApp, deploymentID string) { } func applyScheduledOccurrence(ctx context.Context, database *db.DB, input db.ApplyScheduledOccurrenceInput) error { - return database.Transaction(ctx, func(tx *yourbatis.Tx) error { - return database.ApplyScheduledOccurrenceTx(ctx, tx, input) - }) + return database.ApplyScheduledOccurrence(ctx, input) } func startDeploymentScheduler(t *testing.T, app *testApp) func() { From 2793fb1046d296f7f90dc99d1d79f67a3b0850d4 Mon Sep 17 00:00:00 2001 From: xgxgx Date: Wed, 12 Aug 2026 08:30:42 +0800 Subject: [PATCH 10/10] refactor(agents): remove unrelated handler diff --- internal/agents/handler.go | 9 +++------ 1 file changed, 3 insertions(+), 6 deletions(-) diff --git a/internal/agents/handler.go b/internal/agents/handler.go index 435e07d2..1c8aa80a 100644 --- a/internal/agents/handler.go +++ b/internal/agents/handler.go @@ -100,10 +100,7 @@ type agentReference struct { func NewHandler(cfg config.Config, database *db.DB, logger *slog.Logger) *Handler { logger = logging.LoggerOrDefault(logger) - h := &Handler{ - cfg: cfg, db: database, - errorAdapter: httpapi.NewErrorAdapter(logger), - } + h := &Handler{cfg: cfg, db: database, errorAdapter: httpapi.NewErrorAdapter(logger)} wrap := h.errorAdapter.Wrap router := chi.NewRouter() router.NotFound(wrap(h.notFound)) @@ -370,14 +367,14 @@ func (h *Handler) archive(w http.ResponseWriter, r *http.Request, agentID string httpapi.WriteJSON(w, http.StatusOK, h.fixtureAgent(agentID, 1, true)) return nil } - archived, err := h.db.ArchiveAgent(r.Context(), principal.WorkspaceUUID, agentID) + record, err := h.db.ArchiveAgent(r.Context(), principal.WorkspaceUUID, agentID) if err != nil { if errors.Is(err, db.ErrNotFound) { return agentNotFound(agentID, err) } return internalError("Could not archive agent", fmt.Errorf("archive agent %q: %w", agentID, err)) } - httpapi.WriteJSON(w, http.StatusOK, responseFromAgent(archived)) + httpapi.WriteJSON(w, http.StatusOK, responseFromAgent(record)) return nil }