Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions cmd/migrate/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down Expand Up @@ -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
}
50 changes: 49 additions & 1 deletion docs/design/be/deployments-api-contract.md
Original file line number Diff line number Diff line change
Expand Up @@ -51,9 +51,57 @@ 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 统一由 `github.com/robfig/cron/v3` 解析和计算:

- Deployment schedule 由 `robfig/cron/v3` 解析五段 Cron 和 timezone;解析失败时返回参数错误。
- `upcoming_runs_at` 返回最多五个名义 UTC 时刻,不再使用 366 天扫描上限,因此闰日计划有效。
- spring-forward 不存在的墙上时刻不触发;fall-back 重复的墙上时刻触发两次。
- River Job 直接使用名义 occurrence,不增加私有 jitter 算法。

每个 active 且未归档、schedule 非空的 Deployment 对应一个 River Periodic Job,Periodic Job ID 使用 Deployment ID。Deployment 表是配置真源,只持久化 `schedule`,不保存应用自行推进的下一次游标或额外调度版本。Job 携带注册时的 schedule 快照;worker 读取 Deployment 后以当时的执行配置作为本次 occurrence 快照,最终事务锁行后确认 Deployment 仍为 active、schedule 和执行配置均未变化。

每个应用实例启动时从 Deployment 表加载 Periodic Jobs,并每 10 秒从数据库同步一次 registry。这样所有执行实例最终持有相同配置,进程重启不会丢失 schedule,pause/archive/清空 schedule 会移除 Periodic Job,unpause 或修改 schedule 会重新注册。同步完成前已经投递的 Job 会由 worker 根据当前状态和 schedule 快照跳过。单条确定性的存量 schedule 错误记录后跳过,数据库不可用等全局基础设施错误仍使启动失败。

River 的 leader election 保证只有 leader 根据 Cron 推进并投递 Periodic Job,应用不计算、持久化或插入“下一条 Job”。worker 使用 River Job 的 `scheduled_at` 作为名义 occurrence;暂停或停机期间不补跑历史 occurrence,恢复注册后直接等待 Cron 的下一次。River 开源 Periodic Jobs 的调度状态主要在 leader 内存中,官方不承诺强持久性,leader 切换的极短窗口可能跳过一次 occurrence;需要严格不漏的调度时应采用 River Pro durable periodic jobs,而不是在应用层恢复一套游标链。

```mermaid
sequenceDiagram
participant API as Deployment API
participant AppDB as Deployment config
participant Registry as River Periodic Jobs
participant Leader as River leader
participant Worker as Scheduled worker

API->>AppDB: 提交 schedule 或状态
AppDB-->>Registry: 各实例启动加载并每 10 秒同步配置
Leader->>Registry: 按 Cron 计算下一 occurrence
Leader->>Worker: 投递带名义 scheduled_at 的 Job
Worker->>AppDB: 锁定 Deployment,校验 active 与 schedule 快照
Worker->>AppDB: 同一事务写 Run、Session 与 Deployment 状态
Worker-->>Leader: 返回结果,由 River 完成或重试当前 Job
Note over Leader,Worker: 下一次投递继续由 River Periodic Jobs 推进
```

`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;根 Agent 在触发时已删除也会自动归档,且不生成 Run。
- Workspace 已归档时,最终事务拒绝创建 Session,worker 改为记录 `workspace_archived_error` 失败 Run 并自动暂停 Deployment。
- 其他引用或配置失败生成最终失败 Run。只有公开的 14 类 paused-reason error 会自动暂停;`session_rate_limited_error` 与 `session_creation_rejected_error` 不暂停,并继续下一个 occurrence。
- 数据库或进程级失败交给 River 重试;Run、Session 与 Deployment 状态在同一个 Yourbatis 事务中提交或回滚,当前 River Job 由 Worker 返回结果后交给 River 完成。occurrence 唯一索引保证业务事务提交后发生进程故障时,River 重试不会创建重复 Run。
- paused Deployment 仍允许 manual Run。
- paused 或 archived Deployment 的 `upcoming_runs_at` 为空。

组织级最多保留 1,000 个未归档且 schedule 非空的 Deployment。创建以及从无 schedule 更新为有 schedule 时进行 best-effort 计数检查;并发请求可能短暂越过限制,不额外引入 organization 锁或配额计数器。

主要参考资料:

- <https://platform.claude.com/docs/en/api/beta/deployments>
- <https://platform.claude.com/docs/en/api/beta/deployment_runs>
- <https://platform.claude.com/docs/en/managed-agents/scheduled-deployments>
- <https://riverqueue.com/docs/periodic-jobs>
31 changes: 21 additions & 10 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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.4
github.com/superduck-ai/yourbatis v0.1.5
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
)

Expand All @@ -43,33 +46,41 @@ 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
github.com/jackc/pgpassfile v1.0.0 // indirect
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/oauth2 v0.35.0 // indirect
golang.org/x/mod v0.38.0 // indirect
golang.org/x/oauth2 v0.36.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/api v0.0.0-20260414002931-afd174a4e478 // indirect
google.golang.org/genproto/googleapis/rpc v0.0.0-20260420184626-e10c466a9529 // indirect
google.golang.org/grpc v1.80.0 // indirect
google.golang.org/grpc v1.82.1 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
)
Loading
Loading