|
|
|
@@ -24,6 +24,66 @@ func NewWorkflowAdapter(handler *temporal.Handler) *WorkflowAdapter {
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// HandleStart handles workflow start requests.
|
|
|
|
|
// Expects payload: { "namespace": "default", "workflow_id": "...", "workflow_type": "...", "task_queue": "...", "input": {...} }
|
|
|
|
|
func (wa *WorkflowAdapter) HandleStart(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
wa.forwardToTemporal(w, r)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// HandleDescribe handles workflow describe requests.
|
|
|
|
|
// Expects payload: { "namespace": "default", "workflow_id": "..." }
|
|
|
|
|
func (wa *WorkflowAdapter) HandleDescribe(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
wa.forwardToTemporal(w, r)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// HandleList handles workflow list requests.
|
|
|
|
|
// Expects payload: { "namespace": "default", "query": "..." (optional) }
|
|
|
|
|
func (wa *WorkflowAdapter) HandleList(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
wa.forwardToTemporal(w, r)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// HandleHistory handles workflow history requests.
|
|
|
|
|
// Expects payload: { "namespace": "default", "workflow_id": "..." }
|
|
|
|
|
func (wa *WorkflowAdapter) HandleHistory(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
wa.forwardToTemporal(w, r)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// HandleTerminate handles workflow termination.
|
|
|
|
|
// Expects payload: { "namespace": "default", "workflow_id": "...", "reason": "..." }
|
|
|
|
|
func (wa *WorkflowAdapter) HandleTerminate(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
wa.forwardToTemporal(w, r)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// HandleCancel handles workflow cancellation.
|
|
|
|
|
// Expects payload: { "namespace": "default", "workflow_id": "..." }
|
|
|
|
|
func (wa *WorkflowAdapter) HandleCancel(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
wa.forwardToTemporal(w, r)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// HandleSignal handles workflow signal.
|
|
|
|
|
// Expects payload: { "namespace": "default", "workflow_id": "...", "signal_name": "...", "signal_data": {...} }
|
|
|
|
|
func (wa *WorkflowAdapter) HandleSignal(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
wa.forwardToTemporal(w, r)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// HandleQuery handles workflow query.
|
|
|
|
|
// Expects payload: { "namespace": "default", "workflow_id": "...", "query_type": "...", "query_data": {...} }
|
|
|
|
|
func (wa *WorkflowAdapter) HandleQuery(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
wa.forwardToTemporal(w, r)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// HandleReset handles workflow reset.
|
|
|
|
|
// Expects payload: { "namespace": "default", "workflow_id": "...", "reset_type": "..." }
|
|
|
|
|
func (wa *WorkflowAdapter) HandleReset(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
wa.forwardToTemporal(w, r)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// HandleUpdate handles workflow update.
|
|
|
|
|
// Expects payload: { "namespace": "default", "workflow_id": "...", "update_data": {...} }
|
|
|
|
|
func (wa *WorkflowAdapter) HandleUpdate(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
wa.forwardToTemporal(w, r)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// resourceToAction maps X-Resource names to Temporal action names.
|
|
|
|
|
var resourceToAction = map[string]string{
|
|
|
|
|
"execute": "START_WORKFLOW",
|
|
|
|
@@ -74,8 +134,13 @@ func (wa *WorkflowAdapter) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
|
|
|
|
|
// Inject action into body for temporal handler
|
|
|
|
|
payload["action"] = action
|
|
|
|
|
if _, ok := payload["namespace"]; !ok {
|
|
|
|
|
payload["namespace"] = "default"
|
|
|
|
|
|
|
|
|
|
// namespace is required for all workflow operations
|
|
|
|
|
if ns, ok := payload["namespace"].(string); !ok || ns == "" {
|
|
|
|
|
w.Header().Set("Content-Type", "application/json")
|
|
|
|
|
w.WriteHeader(http.StatusBadRequest)
|
|
|
|
|
fmt.Fprintf(w, `{"error":"namespace is required"}`)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
newBody, _ := json.Marshal(payload)
|
|
|
|
@@ -86,17 +151,51 @@ func (wa *WorkflowAdapter) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
wa.temporalHandler.ServeHTTP(w, r)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// GetWorkflowSpec returns the ServiceAdapter spec for workflow service.
|
|
|
|
|
// forwardToTemporal reads the request body, ensures namespace is specified,
|
|
|
|
|
// and forwards to the temporal handler.
|
|
|
|
|
func (wa *WorkflowAdapter) forwardToTemporal(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
// Read request body
|
|
|
|
|
body, err := io.ReadAll(r.Body)
|
|
|
|
|
if err != nil {
|
|
|
|
|
http.Error(w, fmt.Sprintf("failed to read request body: %v", err), http.StatusBadRequest)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
defer r.Body.Close()
|
|
|
|
|
|
|
|
|
|
// Parse JSON to check for namespace
|
|
|
|
|
var payload map[string]interface{}
|
|
|
|
|
if err := json.Unmarshal(body, &payload); err != nil {
|
|
|
|
|
http.Error(w, fmt.Sprintf("invalid JSON payload: %v", err), http.StatusBadRequest)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Ensure namespace is specified (required for Temporal routing)
|
|
|
|
|
namespace, ok := payload["namespace"].(string)
|
|
|
|
|
if !ok || namespace == "" {
|
|
|
|
|
http.Error(w, `"namespace" field required in payload`, http.StatusBadRequest)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Forward to temporal handler by calling it with the request
|
|
|
|
|
// Restore body for temporal handler
|
|
|
|
|
r.Body = io.NopCloser(bytes.NewReader(body))
|
|
|
|
|
r.ContentLength = int64(len(body))
|
|
|
|
|
|
|
|
|
|
// Call temporal handler
|
|
|
|
|
wa.temporalHandler.ServeHTTP(w, r)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// GetSpec returns the ServiceAdapter spec for workflow service.
|
|
|
|
|
// This defines the available resources and methods.
|
|
|
|
|
func GetWorkflowSpec() *Spec {
|
|
|
|
|
return &Spec{
|
|
|
|
|
ServiceName: "workflow",
|
|
|
|
|
Upstream: Upstream{
|
|
|
|
|
URL: "grpc://temporal:7233",
|
|
|
|
|
URL: "grpc://temporal:7233", // gRPC endpoint
|
|
|
|
|
TimeoutSeconds: 30,
|
|
|
|
|
},
|
|
|
|
|
Auth: Auth{
|
|
|
|
|
Required: true,
|
|
|
|
|
Required: false,
|
|
|
|
|
Capability: "workflow:execute",
|
|
|
|
|
},
|
|
|
|
|
Retryable: true,
|
|
|
|
|