项目文件夹

文件
Asim Aslam 7e2346b8c8 flow: human-in-the-loop pause/resume (durable workflow, stage A) (#4852)
Adds a waiting run state so a flow step can suspend for external input
and resume durably — stage A of the durable-agentic-workflow design in
#4816.

- flow.Await(key, prompt) / flow.AwaitStep(...): a StepFunc that suspends
  the run. runFrom recognizes the signal, checkpoints the run with status
  "waiting" (recording what it awaits), and returns cleanly — a suspend is
  not a failure, and it is not retried or graded.
- Flow.ResumeWith(ctx, runID, input): completes the awaited step with the
  injected input (which becomes that step's output state) and continues
  from the next step.
- Flow.Waiting(ctx): lists suspended runs with their Await metadata.
- ResumePending/Pending skip waiting runs — they need input, not a
  restart. Existing crash-resume (Resume) is unchanged.

Additive: no signature or default-behavior changes. Await ergonomics
(sentinel-return) are the default proposed in #4816; open to AwaitStep-kind
instead if preferred.


Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL

Co-authored-by: Claude <noreply@anthropic.com>
2026-07-15 14:08:27 +01:00

118 行
3.8 KiB
Go

package flow
import (
"context"
"time"
"go-micro.dev/v6/ai"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/codes"
"go.opentelemetry.io/otel/trace"
)
const flowInstrumentationName = "go-micro.dev/v6/flow"
const (
spanNameFlowRun = "flow.run"
spanNameFlowStep = "flow.step"
AttrFlowRunID = "flow.run.id"
AttrFlowParentID = "flow.run.parent_id"
AttrFlowName = "flow.name"
AttrFlowStepName = "flow.step.name"
AttrFlowStatus = "flow.status"
AttrFlowAttempts = "flow.step.attempts"
AttrFlowLatencyMS = "flow.latency_ms"
AttrFlowErrorKind = "flow.error.kind"
AttrFlowVerificationStatus = "flow.verification.status"
AttrFlowVerificationNote = "flow.verification.note"
AttrFlowDispatch = "flow.dispatch"
AttrFlowTrigger = "flow.trigger"
)
func (f *Flow) tracer() trace.Tracer {
return f.opts.TraceProvider.Tracer(flowInstrumentationName)
}
func (f *Flow) startRunSpan(ctx context.Context, run Run) (context.Context, func(Run, error)) {
if f.opts.TraceProvider == nil {
return ctx, func(Run, error) {}
}
info, _ := ai.RunInfoFrom(ctx)
attrs := []attribute.KeyValue{
attribute.String(AttrFlowRunID, run.ID),
attribute.String(AttrFlowParentID, run.ParentID),
attribute.String(AttrFlowName, f.name),
attribute.String(AttrFlowStatus, run.Status),
}
attrs = appendRunInfoDispatch(attrs, info)
ctx, span := f.tracer().Start(ctx, spanNameFlowRun, trace.WithSpanKind(trace.SpanKindInternal), trace.WithAttributes(attrs...))
start := time.Now()
return ctx, func(done Run, err error) {
span.SetAttributes(
attribute.String(AttrFlowStatus, done.Status),
attribute.Int64(AttrFlowLatencyMS, time.Since(start).Milliseconds()),
)
if err != nil {
span.RecordError(err)
span.SetAttributes(attribute.String(AttrFlowErrorKind, string(ai.ClassifyError(err))))
span.SetStatus(codes.Error, err.Error())
} else {
span.SetStatus(codes.Ok, "")
}
span.End()
}
}
func (f *Flow) runStepSpan(ctx context.Context, step Step, in State) (State, int, Verification, error) {
if f.opts.TraceProvider == nil {
return f.runStep(ctx, step, in)
}
info, _ := ai.RunInfoFrom(ctx)
attrs := []attribute.KeyValue{
attribute.String(AttrFlowRunID, info.RunID),
attribute.String(AttrFlowParentID, info.ParentID),
attribute.String(AttrFlowName, f.name),
attribute.String(AttrFlowStepName, step.Name),
}
attrs = appendRunInfoDispatch(attrs, info)
ctx, span := f.tracer().Start(ctx, spanNameFlowStep, trace.WithAttributes(attrs...))
start := time.Now()
out, attempts, verification, err := f.runStep(ctx, step, in)
span.SetAttributes(
attribute.Int(AttrFlowAttempts, attempts),
attribute.Int64(AttrFlowLatencyMS, time.Since(start).Milliseconds()),
)
if verification.Passed {
span.SetAttributes(attribute.String(AttrFlowVerificationStatus, "passed"))
}
if verification.Feedback != "" {
span.SetAttributes(attribute.String(AttrFlowVerificationNote, verification.Feedback))
if !verification.Passed {
span.SetAttributes(attribute.String(AttrFlowVerificationStatus, "failed"))
}
}
if a, ok := isAwaitInput(err); ok {
// A suspend is normal control flow, not a step error.
span.SetStatus(codes.Ok, "waiting: "+a.Key)
} else if err != nil {
span.RecordError(err)
span.SetAttributes(attribute.String(AttrFlowErrorKind, string(ai.ClassifyError(err))))
span.SetStatus(codes.Error, err.Error())
} else {
span.SetStatus(codes.Ok, "")
}
span.End()
return out, attempts, verification, err
}
func appendRunInfoDispatch(attrs []attribute.KeyValue, info ai.RunInfo) []attribute.KeyValue {
if info.Dispatch != "" {
attrs = append(attrs, attribute.String(AttrFlowDispatch, info.Dispatch))
}
if info.Trigger != "" {
attrs = append(attrs, attribute.String(AttrFlowTrigger, info.Trigger))
}
return attrs
}