Skip to content

资源身份与短期授权 ​

问题:长任务中的下载 URL 过期了 ​

把签名 URL 直接当作流程输入,看似省事。任务排队一小时后 URL 过期,重试也只能拿到同一个失效地址。日志里的 URL 还可能泄露访问授权。

资源身份和授权寿命不同。稳定资源引用回答“处理哪一个对象”;短期授权回答“这个执行当前能进行什么操作”;传输地址回答“字节从哪里取得”。

原理:拥有者控制授权,执行者只得到必要权限 ​

业务系统拥有资源与存储位置,调度平台持有不透明引用。执行 Attempt 在需要传输时申请短期 GET/PUT 授权,授权包含操作范围、有效期及必要约束。过期时通过同一资源身份续期。

结果上传也应区分分配、上传和提交:先分配稳定目标,上传字节,确认完整性后把资源标为可用。HTTP 200 不能替代校验和或业务提交。

最小例子:稳定引用与瞬时授权 ​

json
{
  "input": { "resource_ref": "asset:photo-123" },
  "grant": {
    "operation": "GET",
    "expires_at": "2030-01-01T00:05:00Z",
    "temporary_url": "https://storage.example/temporary-access"
  }
}

这是教学结构,不是本项目 wire contract。resource_ref 可保存为执行输入;temporary_url 不应当作长期业务身份或写入普通日志。凭据无法通过该示例推断。

方案比较 ​

平台代理全部字节便于集中控制,但会消耗 Server 带宽、内存和连接。执行端直连对象存储减少中转,需要短期最小权限授权、上传确认和断点恢复。长期凭据复制到每个 Worker 操作简单,泄露影响却更大。

同一资源不同版本必须可区分,否则重试可能处理更新后的文件。稳定引用可以映射不可变版本或在执行时冻结版本,依据业务决定。

真实案例:AccessGrant 与资源拥有者 ​

TDP AccessGrant 应用根据已保存授权、Attempt 和 lease_version 向业务 Resource Provider 申请或续期,校验返回值。App Demo 拥有文件元数据和对象存储;Worker 直接传输文件字节,Server 不充当文件中转。

已核对的实现 · 本地代码快照

资源授权用例 backend · c6dbe05c

server/internal/core/accessgrant/application/service.go · 第 61–158 行
符号:Service.Authorize / issue · 核对日期 2026-10-02
来源与提交版本一致

func (service *Service) Authorize(ctx context.Context, request domain.Authorization) (*domain.Grant, error) {
	if request.AttemptID == "" || request.LeaseVersion < 1 || request.ResourceRef == "" || request.IdempotencyKey == "" ||
		(request.Operation != domain.Download && request.Operation != domain.Upload) ||
		(request.ScopeType == "") != (request.ScopeResourceRef == "") || request.ScopeType != "" && request.Operation != domain.Download {
		return nil, errors.New("invalid resource authorization")
	}
	grantID, err := service.ids.New()
	if err != nil {
		return nil, err
	}
	stored, err := service.repository.ReserveBoundGrant(ctx, request, grantID, service.clock.Now().UTC())
	if err != nil {
		return nil, fmt.Errorf("reserve bound AccessGrant: %w", err)
	}
	return service.issue(ctx, stored, request.AttemptID, request.LeaseVersion, request.Operation, request.IdempotencyKey)
}

func (service *Service) issue(ctx context.Context, stored *domain.StoredGrant, attemptID string, leaseVersion int64, operation domain.Operation, idempotencyKey string) (*domain.Grant, error) {
	secret, err := service.cipher.Open(stored.EncryptedBearer, append([]byte("resource-provider:"), stored.ApplicationAAD...))
	if err != nil {
		return nil, err
	}
	issued, err := service.provider.Issue(ctx, domain.IssueRequest{
		GrantID:        stored.GrantID,
		ApplicationID:  stored.ApplicationID,
		ResourceRef:    stored.ResourceRef,
		AttemptID:      attemptID,
		LeaseVersion:   leaseVersion,
		Operation:      operation,
		IdempotencyKey: idempotencyKey,
		ProviderURL:    stored.ProviderURL,
		BearerToken:    string(secret),
		Scope:          stored.Scope,
	})
	if err != nil {
		return nil, fmt.Errorf("issue AccessGrant: %w", err)
	}
	if err := service.validate(*issued, domain.Renewal{
		GrantID:   stored.GrantID,
		Operation: operation,
	}); err != nil {
		return nil, err
	}
	encoded, err := json.Marshal(issued)
	if err != nil {
		return nil, fmt.Errorf("encode AccessGrant: %w", err)
	}
	ciphertext, err := service.cipher.Seal(encoded, stored.GrantAAD)
	if err != nil {
		return nil, err
	}
	if err := service.repository.Save(ctx, stored.GrantID, ciphertext, issued.ExpiresAt, service.clock.Now().UTC()); err != nil {
		return nil, fmt.Errorf("save AccessGrant: %w", err)
	}
	return issued, nil
}

func (service *Service) Allocate(ctx context.Context, request domain.AllocationRequest) (*domain.Allocation, error) {
	if request.AttemptID == "" || request.LeaseVersion < 1 || request.OutputCollectionRef == "" || request.RelativePath == "" ||
		request.ArtifactType == "" || len(request.ArtifactType) > 100 || request.MediaType == "" || request.MaxBytes < 1 {
		return nil, errors.New("invalid output artifact allocation")
	}
	allocationID, err := service.ids.New()
	if err != nil {
		return nil, err
	}
	stored, err := service.repository.ReserveAllocation(ctx, request, allocationID, service.clock.Now().UTC())
	if err != nil {
		return nil, fmt.Errorf("reserve output artifact allocation: %w", err)
	}
	secret, err := service.cipher.Open(stored.EncryptedSecret, append([]byte("resource-provider:"), stored.ApplicationAAD...))
	if err != nil {
		return nil, err
	}
	allocated, err := service.provider.Allocate(ctx, domain.AllocateProviderRequest{
		StoredAllocation: *stored,
		BearerToken:      string(secret),
	})
	if err != nil {
		return nil, fmt.Errorf("allocate output artifact: %w", err)
	}
	if allocated.ResourceRef == "" || allocated.Grant.Operation != domain.Upload || allocated.Grant.Method != "PUT" || allocated.Grant.MaxBytes < 1 ||
		allocated.Grant.ID == "" || !allocated.Grant.ExpiresAt.After(service.clock.Now()) || !validTemporaryURL(allocated.Grant.URL, service.allowInsecureHTTP) {
		return nil, errors.New("Resource Provider returned an invalid output allocation")
	}
	encoded, err := json.Marshal(&allocated.Grant)
	if err != nil {
		return nil, err
	}
	grantAAD := []byte(allocated.Grant.ID)
	ciphertext, err := service.cipher.Seal(encoded, grantAAD)
	if err != nil {
		return nil, err
	}
	if err := service.repository.CompleteAllocation(ctx, *stored, *allocated, ciphertext, service.clock.Now().UTC()); err != nil {
		return nil, fmt.Errorf("complete output artifact allocation: %w", err)
	}
	return allocated, nil

片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。

跨项目端到端资料 workbench · 514f5a00

docs/learning/README.md · 第 11–57 行
符号:主链路与事实边界 · 核对日期 2026-10-02
来源与提交版本一致

## 最有效的学习方式

不要从目录树第一行开始顺序读完整个项目。先跟一条真实的 `ARCHIVE_CLEAN_THEN_WATERMARK` 处理单纵向走完,再回头学习每个模块的横向设计。

这条纵向链路覆盖了项目最重要的机制:

```text
浏览器上传 ZIP
  → app-demo 创建 ProcessingOrder 和本地 Command Outbox
  → app-demo 后台提交 TDP Task
  → TDP 创建 Task / Execution / Event / Outbox
  → NATS JetStream 唤醒 Temporal
  → Temporal 展开 DAG 并创建 Step / Attempt
  → PostgreSQL Command 经 SSE 到 Worker
  → Worker SQLite Inbox → archive.clean → 本地 Event Outbox
  → Server Attempt Inbox → Temporal Signal
  → archive.clean 输出经 CEL 映射给 image.watermark
  → Execution 完成并可靠 Callback
  → app-demo Callback Inbox 更新业务投影
```

完整细节见[一条处理单的数据流](business-flow.md)。

## 先建立四个“事实边界”

| 边界 | 保存内容 | 是否可由其他系统直接改写 |
|---|---|---|
| app-demo PostgreSQL | 处理单、业务状态、资源元数据、提交 Outbox、Callback Inbox | 否;TDP 只能通过 API/Callback 影响业务投影 |
| TDP PostgreSQL | Task、Execution、Step、Attempt、Worker、租约、事件和可靠消息 | 否;Temporal、NATS 和 Worker 都不能绕过应用用例直接成为事实源 |
| Worker SQLite | 当前 Worker 的命令 Inbox、Attempt、本地 Checkpoint、事件 Outbox、注册身份 | 否;只属于该 Worker,不是全局查询库 |
| MinIO | app-demo 拥有的 ZIP、PNG 和 manifest 文件字节 | 通过短期 GET/PUT 授权访问;TDP Server 不转发文件内容 |

Temporal 保存耐久编排历史,NATS 保存可靠传输中的消息;它们都不替代 PostgreSQL 查询模型。删除 Worker SQLite 也不是“清缓存”,而是放弃该实例尚未同步的恢复状态。

## 文档地图

1. [一条处理单的数据流](business-flow.md)
   - 以 `ARCHIVE_CLEAN_THEN_WATERMARK` 为主线,从上传、创建、调度、Worker 执行一直跟到 Callback。
   - 重点看数据在业务 ID、资源引用、AccessGrant、Step 输出和业务投影之间如何变形。
2. [数据表与一致性边界](data-model.md)
   - 分别解释 app-demo PostgreSQL、TDP PostgreSQL、Worker SQLite 的表设计。
   - 不只列字段,还说明每组表为何要在同一个事务里写、哪些列是事实、哪些列是投影。
3. [代码阅读与联调练习](code-reading.md)
   - 给出按业务动作定位代码的方法,以及一组不会删除本地数据的观察命令。
   - 用同一组 ID 串起页面、日志、两套 PostgreSQL、Temporal UI 和 NATS NUI。

已有文档仍然是重要的专题参考:

片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。

ZIP 清理、水印和缩略图把资源契约转为具体文件处理。通用运行时管理执行与事件,能力负责格式、安全限制与产物语义。

失败与边界:续期并非无条件放行 ​

资源已删除、Attempt 失效、租约版本改变或权限撤销时,续期应失败。执行端不能因暂时访问失败而自行扩大权限。

输入路径需要防 SSRF、任意本地路径访问和压缩包穿越。预签名 URL 也可能包含权限,日志仅记录资源引用与安全元数据。中断上传后应确认目标状态,避免把半成品标成完成。

迁移练习与参考答案 ​

练习:模型训练任务排队很久,数据集 URL 常过期。怎样改变接口?

参考答案:输入传稳定且版本明确的数据集引用;训练开始或下载重试时由数据集拥有者签发最小权限短期授权。保存下载 Checkpoint 和内容校验,续期前验证任务资格。训练输出分配独立资源身份并在完整上传后提交。

继续阅读:身份与安全边界、Worker 恢复。

理解原理 · 分析取舍 · 用真实代码检验