| 3740 | } |
| 3741 | |
| 3742 | func (m *Service) recordTaskEvent( |
| 3743 | ctx context.Context, |
| 3744 | taskID string, |
| 3745 | runID string, |
| 3746 | eventType string, |
| 3747 | actor ActorContext, |
| 3748 | payload any, |
| 3749 | ) error { |
| 3750 | rawPayload, err := marshalTaskEventPayload(payload) |
| 3751 | if err != nil { |
| 3752 | return err |
| 3753 | } |
| 3754 | event := Event{ |
| 3755 | ID: m.newID("evt"), |
| 3756 | TaskID: strings.TrimSpace(taskID), |
| 3757 | RunID: strings.TrimSpace(runID), |
| 3758 | EventType: strings.TrimSpace(eventType), |
| 3759 | Actor: actor.Actor, |
| 3760 | Origin: actor.Origin, |
| 3761 | Payload: rawPayload, |
| 3762 | Timestamp: m.now().UTC(), |
| 3763 | } |
| 3764 | if err := m.store.CreateTaskEvent(ctx, event); err != nil { |
| 3765 | return err |
| 3766 | } |
| 3767 | |
| 3768 | postCommitCtx := context.Background() |
| 3769 | if ctx != nil { |
| 3770 | postCommitCtx = context.WithoutCancel(ctx) |
| 3771 | } |
| 3772 | |
| 3773 | record, err := m.store.GetTaskEventRecord(postCommitCtx, event.ID) |
| 3774 | if err != nil { |
| 3775 | m.emitTaskLiveEventBestEffort(postCommitCtx, event.ID) |
| 3776 | return nil |
| 3777 | } |
| 3778 | m.notifyTaskObserverBestEffort(postCommitCtx, record) |
| 3779 | m.emitTaskLiveRecordBestEffort(postCommitCtx, record) |
| 3780 | return nil |
| 3781 | } |
| 3782 | |
| 3783 | func (m *Service) notifyTaskObserverBestEffort(ctx context.Context, record EventRecord) { |
| 3784 | if m == nil || m.eventObserver == nil { |