Compare commits
14
Commits
main
..
ef79bf9a80
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ef79bf9a80 | ||
|
|
db8f103fa5 | ||
|
|
234d2a2c02 | ||
|
|
d16597ca0b | ||
|
|
509cb39521 | ||
|
|
03f3239cd0 | ||
|
|
d5c7440655 | ||
|
|
bff46fefe7 | ||
|
|
c93bfca6a8 | ||
|
|
44ec3502dc | ||
|
|
6608f1a8d5 | ||
|
|
c6cd41fd4b | ||
|
|
b0f608f145 | ||
|
|
cc9a32f53a |
@@ -0,0 +1,4 @@
|
|||||||
|
(apply,CacheStats{hitCount=337, missCount=199, loadSuccessCount=199, loadExceptionCount=0, totalLoadTime=581291927, evictionCount=0})
|
||||||
|
(tree,CacheStats{hitCount=986, missCount=352, loadSuccessCount=299, loadExceptionCount=0, totalLoadTime=821650758, evictionCount=0})
|
||||||
|
(commit,CacheStats{hitCount=108, missCount=107, loadSuccessCount=107, loadExceptionCount=0, totalLoadTime=78983052, evictionCount=0})
|
||||||
|
(tag,CacheStats{hitCount=0, missCount=2, loadSuccessCount=2, loadExceptionCount=0, totalLoadTime=319542, evictionCount=0})
|
||||||
@@ -0,0 +1,4 @@
|
|||||||
|
e71e5b78236a67327c678490cb50b46981f19de0 bbcbb68b91e786eb71bbb0a4443d7b8a26140e1b .sops.yaml
|
||||||
|
4189696f5581ac0ffdc125c3bf9b9f664b3ddfb0 7cd3f1ee4865c563d141464f6fc185436993b84b .sops.yaml
|
||||||
|
635630e73152a5f22e6cbd42322ec55d79f8d9c0 297e94a89d73d18c4f47013bb0e8303f123715f3 configmap.yaml
|
||||||
|
29e515e7b46742fab8c3fcc2189af7010a6ccc62 6869fa11f96e03f7ec76a0ea14a4ddaf604004a4 gateway-config-secret.enc.yaml
|
||||||
@@ -0,0 +1,12 @@
|
|||||||
|
0a95af80c0051bacbeb8483c1632e47acd3db5be 40207e487cfb63409a976fb2a0b9e1e62c8b1513
|
||||||
|
27428d910111299d0699f429190284a9ca6e50b7 3318daf758349402aef43b095482743ab96b37f9
|
||||||
|
329a495af4c935529fdae17229314101c0c77876 67f24ea76359c8dba4b56267790aad76bbc58464
|
||||||
|
4c8bc6c920b6b75399555827022f69ef0c4f7d15 1fa839b41975fa3f0ac9052355ffb625f5a8f324
|
||||||
|
528545f414c83217408edfea234dcd1f3edee0c2 b8f95506ca1545b876b5531cd385172e9ca5b4b0
|
||||||
|
81038e1cf7567a9133d7c233a97b1e2f19fa1c82 4a00312906ba725f3968187656fde2663b1763ab
|
||||||
|
a5b3b5c44a406896bcb414df6c6426c277715706 2ab47a9dbe5ba36dfa0e275991ef7b7656908410
|
||||||
|
ce27643667a0399115cd1f2b6d38123fdcf2b4f1 6ff0a50de8efbad105fa588245f22fdb26afddc4
|
||||||
|
d49756886a46542b38533b913a1f776b5145f5ec d82cc5a6970a1fb32e21dda9a737b987a8668111
|
||||||
|
db3a30fbcf1f139c667fb68a91762582c49b8cee 04619a269fed9eeea53ab4d4d73131e3713f40a0
|
||||||
|
eb54715e4dec0fb35402576fcc224a09808b00c1 d53b7632cf9646dda1c978a5f94615dc9eaed5e8
|
||||||
|
ef72b5bbccf2df89aa1c86dee29311c63f33bf62 ba55d184fefef1a73a50409ca4fb1f7b27f5b075
|
||||||
+1
-5
@@ -108,12 +108,8 @@ func main() {
|
|||||||
}
|
}
|
||||||
_ = registry.Add(notifAdapter)
|
_ = registry.Add(notifAdapter)
|
||||||
|
|
||||||
// Add other adapters from config (skip if already registered in code)
|
// Add other adapters from config
|
||||||
for _, a := range cfg.Adapters {
|
for _, a := range cfg.Adapters {
|
||||||
if existing := registry.Get(a.ServiceName); existing != nil {
|
|
||||||
log.Printf("skip config adapter '%s': already registered with internal handler", a.ServiceName)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
_ = registry.Add(a)
|
_ = registry.Add(a)
|
||||||
}
|
}
|
||||||
log.Printf("%d service adapters loaded", registry.Count())
|
log.Printf("%d service adapters loaded", registry.Count())
|
||||||
|
|||||||
@@ -24,66 +24,6 @@ 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.
|
// resourceToAction maps X-Resource names to Temporal action names.
|
||||||
var resourceToAction = map[string]string{
|
var resourceToAction = map[string]string{
|
||||||
"execute": "START_WORKFLOW",
|
"execute": "START_WORKFLOW",
|
||||||
@@ -134,13 +74,8 @@ func (wa *WorkflowAdapter) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||||||
|
|
||||||
// Inject action into body for temporal handler
|
// Inject action into body for temporal handler
|
||||||
payload["action"] = action
|
payload["action"] = action
|
||||||
|
if _, ok := payload["namespace"]; !ok {
|
||||||
// namespace is required for all workflow operations
|
payload["namespace"] = "default"
|
||||||
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)
|
newBody, _ := json.Marshal(payload)
|
||||||
@@ -151,51 +86,17 @@ func (wa *WorkflowAdapter) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||||||
wa.temporalHandler.ServeHTTP(w, r)
|
wa.temporalHandler.ServeHTTP(w, r)
|
||||||
}
|
}
|
||||||
|
|
||||||
// forwardToTemporal reads the request body, ensures namespace is specified,
|
// GetWorkflowSpec returns the ServiceAdapter spec for workflow service.
|
||||||
// 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.
|
// This defines the available resources and methods.
|
||||||
func GetWorkflowSpec() *Spec {
|
func GetWorkflowSpec() *Spec {
|
||||||
return &Spec{
|
return &Spec{
|
||||||
ServiceName: "workflow",
|
ServiceName: "workflow",
|
||||||
Upstream: Upstream{
|
Upstream: Upstream{
|
||||||
URL: "grpc://temporal:7233", // gRPC endpoint
|
URL: "grpc://temporal:7233",
|
||||||
TimeoutSeconds: 30,
|
TimeoutSeconds: 30,
|
||||||
},
|
},
|
||||||
Auth: Auth{
|
Auth: Auth{
|
||||||
Required: false,
|
Required: true,
|
||||||
Capability: "workflow:execute",
|
Capability: "workflow:execute",
|
||||||
},
|
},
|
||||||
Retryable: true,
|
Retryable: true,
|
||||||
|
|||||||
Reference in New Issue
Block a user