MCPcopy Create free account
hub / github.com/compozy/agh / EnqueueRunFromPeer

Method EnqueueRunFromPeer

internal/network/tasks.go:292–343  ·  view source on GitHub ↗

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,
)

Source from the content-addressed store, hash-verified

290// EnqueueRunFromPeer enqueues one task run from an authenticated network peer
291// while preserving origin-scoped idempotency inside the task manager.
292func (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
345func (m *Manager) resolveTaskPeerContext(
346 ctx context.Context,

Calls 8

rejectTaskIngressMethod · 0.95
recordTaskIngressMethod · 0.95
enforceBoundTaskChannelFunction · 0.85
withNetworkRunMetadataFunction · 0.85
GetTaskMethod · 0.65
EnqueueRunMethod · 0.65