EnqueueRunFromPeer enqueues one task run from an authenticated network peer while preserving origin-scoped idempotency inside the task manager.
( ctx context.Context, ingress TaskIngressContext, spec taskpkg.EnqueueRun, )
| 290 | // EnqueueRunFromPeer enqueues one task run from an authenticated network peer |
| 291 | // while preserving origin-scoped idempotency inside the task manager. |
| 292 | func (m *Manager) EnqueueRunFromPeer( |
| 293 | ctx context.Context, |
| 294 | ingress TaskIngressContext, |
| 295 | spec taskpkg.EnqueueRun, |
| 296 | ) (*taskpkg.Run, error) { |
| 297 | peerCtx, err := m.resolveTaskPeerContext(ctx, ingress, networkTaskActionEnqueue) |
| 298 | if err != nil { |
| 299 | return nil, err |
| 300 | } |
| 301 | view, err := m.tasks.GetTask(ctx, strings.TrimSpace(spec.TaskID), peerCtx.actor) |
| 302 | if err != nil { |
| 303 | return nil, m.rejectTaskIngress(ctx, peerCtx.ingress, networkTaskActionEnqueue, err, nil) |
| 304 | } |
| 305 | if err := enforceBoundTaskChannel( |
| 306 | view.Task.ID, |
| 307 | view.Task.NetworkChannel, |
| 308 | peerCtx.ingress.Channel, |
| 309 | nil, |
| 310 | ); err != nil { |
| 311 | return nil, m.rejectTaskIngress(ctx, peerCtx.ingress, networkTaskActionEnqueue, err, map[string]any{ |
| 312 | tasksTaskIDKey: view.Task.ID, |
| 313 | tasksNetworkChannelKey: strings.TrimSpace(view.Task.NetworkChannel), |
| 314 | }) |
| 315 | } |
| 316 | if err := validateRequestedTaskChannel(peerCtx.ingress.Channel, spec.NetworkChannel); err != nil { |
| 317 | return nil, m.rejectTaskIngress(ctx, peerCtx.ingress, networkTaskActionEnqueue, err, map[string]any{ |
| 318 | tasksTaskIDKey: view.Task.ID, |
| 319 | tasksNetworkChannelKey: strings.TrimSpace(spec.NetworkChannel), |
| 320 | }) |
| 321 | } |
| 322 | spec, err = withNetworkRunMetadata(spec, peerCtx.ingress) |
| 323 | if err != nil { |
| 324 | return nil, m.rejectTaskIngress(ctx, peerCtx.ingress, networkTaskActionEnqueue, err, map[string]any{ |
| 325 | tasksTaskIDKey: view.Task.ID, |
| 326 | }) |
| 327 | } |
| 328 | |
| 329 | run, err := m.tasks.EnqueueRun(ctx, spec, peerCtx.actor) |
| 330 | if err != nil { |
| 331 | return nil, m.rejectTaskIngress(ctx, peerCtx.ingress, networkTaskActionEnqueue, err, map[string]any{ |
| 332 | tasksTaskIDKey: view.Task.ID, |
| 333 | "idempotency_key": strings.TrimSpace(spec.IdempotencyKey), |
| 334 | }) |
| 335 | } |
| 336 | m.recordTaskIngress(ctx, peerCtx.ingress, networkTaskActionEnqueue, AuditDirectionReceived, "", map[string]any{ |
| 337 | tasksTaskIDKey: run.TaskID, |
| 338 | "run_id": run.ID, |
| 339 | "idempotency_key": strings.TrimSpace(run.IdempotencyKey), |
| 340 | tasksNetworkChannelKey: strings.TrimSpace(run.NetworkChannel), |
| 341 | }) |
| 342 | return run, nil |
| 343 | } |
| 344 | |
| 345 | func (m *Manager) resolveTaskPeerContext( |
| 346 | ctx context.Context, |