nx_wflow_result.nx source
↩ module page · 291 lines · 17179 B
1// Additive workflow result adapter. Existing fire/resume remain unchanged.
2// Whole diagnostics are retained; ledger completion is not served acceptance.
3import "nx_wflow_engine.nx"
4import "nx_sha256.nx"
5import "json_emit.nx"
6struct WfrFile {text:*WfOwnedText,status:i64,hash:*u8}
7struct WfrRun {rid:*u8,state:i64,oks:i64,attempts:i64,failedstep:i64,laststep:i64,pending:i64}
8func wfr_read(path:*u8)->*WfrFile{
9 let f:*WfrFile=sys_mmap(__size_of(WfrFile)) as *WfrFile
10 f.status=0-1
11 let fd:i64=sys_openat_rd(path)
12 if fd<0{if fd==(0-2){f.status=0-2};return f}
13 let n:i64=sys_lseek(fd,0,2)
14 if n<0 || n==WF_TEXT_I64_MAX{sys_close(fd);return f}
15 if sys_lseek(fd,0,0)!=0{sys_close(fd);return f}
16 f.text=wf_text_new(n,1)
17 if f.text==(0 as *WfOwnedText){sys_close(fd);return f}
18 var at:i64=0
19 while at<n{
20 let got:i64=sys_read(fd,((f.text.data as i64)+at) as *u8,n-at)
21 if got==(0-4){}else{if got<=0{sys_close(fd);f.status=0-3;return f};at=at+got}
22 }
23 let probe:*u8=sys_mmap(1)
24 var tail:i64=sys_read(fd,probe,1)
25 while tail==(0-4){tail=sys_read(fd,probe,1)}
26 let end:i64=sys_lseek(fd,0,2);sys_close(fd)
27 if tail!=0 || end!=n{f.status=0-4;return f}
28 let digest:*u8=sys_mmap(SHA256_DIGEST_BYTES)
29 if sha256_digest_checked_native(f.text.data,n,digest)!=0{return f}
30 f.hash=sys_mmap(SHA256_DIGEST_BYTES*2+1)
31 let hex:*u8="0123456789abcdef"
32 var k:i64=0;while k<SHA256_DIGEST_BYTES{f.hash[k*2]=hex[(digest[k] as i64)>>4];f.hash[k*2+1]=hex[(digest[k] as i64)&15];k=k+1}
33 f.status=1;return f
34}
35func wfr_file_state(f:*WfrFile)->*u8{
36 if f.status==1{return "complete"};if f.status==(0-2){return "missing"}
37 if f.status==(0-3){return "short-read"};if f.status==(0-4){return "changed-during-read"};return "unreadable"
38}
39func wfr_frame(buf:*u8,n:i64)->i64{
40 if n==0{return 0};if buf[n-1]!=(10 as u8){return 0-2}
41 var s:i64=0
42 while s<n{
43 var e:i64=s;while e<n{if buf[e]==(0 as u8){return 0-1};if buf[e]==(10 as u8){break};e=e+1}
44 var known:i64=0
45 if wf_text_match(buf,s,e,"WFRUN rid=",slen("WFRUN rid="))==1{known=1}
46 if wf_text_match(buf,s,e,"WFVAR rid=",slen("WFVAR rid="))==1{known=1}
47 if wf_text_match(buf,s,e,"WFVER rid=",slen("WFVER rid="))==1{known=1}
48 if wf_text_match(buf,s,e,"WFDEC rid=",slen("WFDEC rid="))==1{known=1}
49 if wf_text_match(buf,s,e,"WFVERDRIFT rid=",slen("WFVERDRIFT rid="))==1{known=1}
50 if known==0{return 0-1};s=e+1
51 };return 1
52}
53func wfr_field(buf:*u8,s:i64,e:i64,key:*u8)->*WfOwnedText{
54 let t:*WfOwnedText=wf_text_new(e-s,1)
55 if t==(0 as *WfOwnedText){return t}
56 let n:i64=wf_val_after(buf,s,e,key,t.data,t.bytes)
57 if n<0{wf_text_free(t);return 0 as *WfOwnedText}
58 t.len=slen(t.data);return t
59}
60func wfr_scan(buf:*u8,n:i64,rid:*u8,out:*WfrRun)->i64{
61 out.rid=rid;out.state=0;out.pending=0;var s:i64=0;var starts:i64=0;var terminal:i64=0
62 while s<n{
63 var e:i64=s;while e<n{if buf[e]==(10 as u8){break};e=e+1}
64 if wf_text_match(buf,s,e,"WFRUN rid=",slen("WFRUN rid="))==1{
65 let id:*WfOwnedText=wfr_field(buf,s,e," rid=")
66 if id==(0 as *WfOwnedText){return 0-1}
67 if seq(id.data,rid)==1{
68 let st:*WfOwnedText=wfr_field(buf,s,e," status=")
69 let step:*WfOwnedText=wfr_field(buf,s,e," step=")
70 if st==(0 as *WfOwnedText) || step==(0 as *WfOwnedText){return 0-1}
71 let idx:i64=wf_atoi(step.data);if idx<0{return 0-1};out.laststep=idx
72 var known:i64=0
73 if seq(st.data,"START")==1{known=1;starts=starts+1;if starts!=1{return 0-1}}
74 if seq(st.data,"BOUND")==1{known=1}
75 if seq(st.data,"ATT")==1{known=1;out.attempts=out.attempts+1;out.pending=idx;out.state=0}
76 if seq(st.data,"OK")==1{known=1;out.oks=out.oks+1;if out.pending==idx{out.pending=0};out.state=0}
77 if seq(st.data,"FAILSTEP")==1{known=1;out.failedstep=idx;if out.pending==idx{out.pending=0}}
78 if seq(st.data,"WAIT")==1{known=1;out.state=1}
79 if seq(st.data,"FAILED")==1{known=1;if terminal!=0{return 0-1};terminal=2;out.state=2}
80 if seq(st.data,"DONE")==1{known=1;if terminal!=0{return 0-1};terminal=3;out.state=3}
81 if known==0{return 0-1}
82 if terminal!=0{if out.state!=terminal{return 0-1}}
83 wf_text_free(st);wf_text_free(step)
84 };wf_text_free(id)
85 };s=e+1
86 }
87 if starts!=1{return 0-1}
88 if out.state==3 && out.pending!=0{return 0-1}
89 return 0
90}
91func wfr_js(w:*JsonWriter,key:*u8,val:*u8)->i64{
92 if json_emit_key(w,key,slen(key))<0{return 0-1};return json_emit_string(w,val,slen(val))
93}
94func wfr_ji(w:*JsonWriter,key:*u8,val:i64)->i64{
95 if json_emit_key(w,key,slen(key))<0{return 0-1};return json_emit_int(w,val)
96}
97func wfr_file_json(w:*JsonWriter,key:*u8,path:*u8,f:*WfrFile)->i64{
98 if json_emit_key(w,key,slen(key))<0{return 0-1};if json_begin_object(w)<0{return 0-1}
99 if wfr_js(w,"path",path)<0{return 0-1};if wfr_js(w,"state",wfr_file_state(f))<0{return 0-1}
100 if f.status==1{if wfr_ji(w,"bytes",f.text.len)<0{return 0-1};if wfr_js(w,"sha256",f.hash)<0{return 0-1}}
101 return json_end_object(w)
102}
103func wfr_write_all(fd:i64,buf:*u8,n:i64)->i64{
104 var at:i64=0;while at<n{let r:i64=sys_write(fd,((buf as i64)+at) as *u8,n-at);if r==(0-4){}else{if r<=0{return 0-1};at=at+r}};return 0
105}
106func wfr_match_id(rid:*u8,filter:*u8,mode:i64)->i64{
107 if mode==2{var s:i64=0;let n:i64=slen(filter);while s<n{var e:i64=s;while e<n{if filter[e]==(10 as u8){break};e=e+1};if e-s==slen(rid){if wf_text_match(filter,s,e,rid,e-s)==1{return 1}};s=e+1};return 0}
108 if seq(filter,"*")==1{return 1}
109 if mode==0{return seq(rid,filter)}
110 let n:i64=slen(filter);let rn:i64=slen(rid)
111 if rn<=n{return 0};if wf_text_match(rid,0,rn,filter,n)!=1{return 0};return rid[n]==(46 as u8)
112}
113func wfr_var_json(w:*JsonWriter,buf:*u8,n:i64,rid:*u8,key:*u8)->i64{
114 let span:*WfVarSpan=sys_mmap(__size_of(WfVarSpan)) as *WfVarSpan
115 if wf_var_span(buf,n,rid,key,slen(key),span)!=1{return 0}
116 if json_emit_key(w,key,slen(key))<0{return 0-1}
117 return json_emit_string(w,((buf as i64)+span.start) as *u8,span.len)
118}
119
120func wfr_run_labels(w:*JsonWriter,buf:*u8,n:i64,rid:*u8)->i64{
121 var err:i64=0;var ls:i64=0;var flow:*WfOwnedText=0 as *WfOwnedText;var version:i64=0
122 while ls<n{var le:i64=ls;while le<n{if buf[le]==(10 as u8){break};le=le+1}
123 if wf_text_match(buf,ls,le,"WFRUN rid=",10)==1{let r:*WfOwnedText=wf_id_field(buf,ls,le," rid=");let st:*WfOwnedText=wf_id_field(buf,ls,le," status=");if r==(0 as *WfOwnedText)||st==(0 as *WfOwnedText){return 0-1}
124 if seq(r.data,rid)==1&&seq(st.data,"START")==1{flow=wf_id_field(buf,ls,le," flow=");let v:*WfOwnedText=wf_id_field(buf,ls,le," identity=");if v==(0 as *WfOwnedText){return 0-1};if seq(v.data,"2")==1{version=2};wf_text_free(v);wf_text_free(r);wf_text_free(st);break};wf_text_free(r);wf_text_free(st)
125 };ls=le+1
126 }
127 err=err+wfr_ji(w,"identity_version",version)
128 if flow!=(0 as *WfOwnedText){err=err+wfr_js(w,"flow_id",flow.data)}
129 let span:*WfVarSpan=sys_mmap_try(__size_of(WfVarSpan)) as *WfVarSpan;if (span as i64)<=0{return 0-1}
130 if wf_var_span(buf,n,rid,"id",2,span)==1{err=err+json_emit_key(w,"event_id",8)+json_emit_string(w,((buf as i64)+span.start) as *u8,span.len)}
131 var binding:*u8="unverified";if flow!=(0 as *WfOwnedText){if wf_id_replay_binding(buf,n,rid,flow.data)==1{binding="verified"}}
132 err=err+wfr_js(w,"identity_binding",binding);wf_text_free(flow);sys_munmap_direct(span as *u8,__size_of(WfVarSpan));return err
133}
134
135func wfr_summary(ledger:*u8,evidence:*u8,filter:*u8,mode:i64,reaped:i64,exitcode:i64)->i64{
136 let f:*WfrFile=wfr_read(ledger);let log:*WfrFile=wfr_read(evidence)
137 var n:i64=0;var frame:i64=0;var journal:*u8=wfr_file_state(f)
138 if f.status==1{n=f.text.len;frame=wfr_frame(f.text.data,n);journal="complete";if frame==0{journal="empty"};if frame==(0-2){journal="truncated-record"};if frame==(0-1){journal="malformed-record"}}
139 var count:i64=0;var s:i64=0
140 if frame==1{while s<n{if f.text.data[s]==(10 as u8){count=count+1};s=s+1}}
141 if count>=WF_TEXT_I64_MAX/__size_of(WfrRun){return 2}
142 let runs:*WfrRun=sys_mmap((count+1)*__size_of(WfrRun)) as *WfrRun
143 var nr:i64=0;s=0
144 if frame==1{
145 while s<n{
146 var e:i64=s;while e<n{if f.text.data[e]==(10 as u8){break};e=e+1}
147 if wf_text_match(f.text.data,s,e,"WFRUN rid=",slen("WFRUN rid="))==1{
148 let st:*WfOwnedText=wfr_field(f.text.data,s,e," status=")
149 if st==(0 as *WfOwnedText){return 2}
150 if seq(st.data,"START")==1{
151 let id:*WfOwnedText=wfr_field(f.text.data,s,e," rid=");if id==(0 as *WfOwnedText){return 2}
152 var selected:i64=wfr_match_id(id.data,filter,mode)
153 if mode==1{selected=1}
154 if selected==1 && mode==1{
155 let identity:*WfVarSpan=sys_mmap(__size_of(WfVarSpan)) as *WfVarSpan
156 selected=0
157 if wf_var_span(f.text.data,n,id.data,"id",2,identity)==1{if identity.len==slen(filter){selected=wf_text_match(f.text.data,identity.start,identity.start+identity.len,filter,identity.len)}}
158 }
159 if selected==1{
160 if wfr_scan(f.text.data,n,id.data,((runs as i64)+nr*__size_of(WfrRun)) as *WfrRun)!=0{journal="malformed-record";frame=0-1}
161 nr=nr+1
162 }else{wf_text_free(id)}
163 };wf_text_free(st)
164 };s=e+1
165 }
166 }
167 // Six output bytes is the RFC8259 worst case for one escaped input byte.
168 // Repeated key/constant overhead follows the measured schema text per run.
169 let scaffold:*u8="schema nishi-workflow-result/1 owner nx_wflow journal_state command_execution_reaped command_exit_code effects connector-receipts-required served_acceptance unverified evidence ledger path state bytes sha256 runs run_id event_id flow_id identity_version identity_binding verified unverified state DONE FAILED PARKED UNCERTAIN ok_steps attempts failed_step pending_step delivery_report manifest manifest_sha256 release_scope delivery_state action inspect-complete-evidence-and-reconcile-connector-receipts do-not-auto-retry failed_steps done failed parked uncertain selected_runs workflow-complete-review-effects-and-served-output malformed-record changed-during-read"
170 let unit:i64=slen(scaffold)
171 if nr>WF_TEXT_I64_MAX/unit-1{return 2}
172 let overhead:i64=(nr+1)*unit
173 if n>WF_TEXT_I64_MAX-overhead{return 2};var needed:i64=n+overhead
174 let paths:i64=slen(ledger)+slen(evidence)+slen(filter)
175 if needed>WF_TEXT_I64_MAX-paths{return 2};needed=needed+paths
176 if needed>WF_TEXT_I64_MAX/6{return 2};needed=needed*6
177 let output:*WfOwnedText=wf_text_new(needed,1);if output==(0 as *WfOwnedText){return 2}
178 let w:*JsonWriter=sys_mmap(__size_of(JsonWriter)) as *JsonWriter;json_writer_init(w,output.data,needed)
179 var err:i64=json_begin_object(w)
180 err=err+wfr_js(w,"schema","nishi-workflow-result/1")+wfr_js(w,"owner","nx_wflow")+wfr_js(w,"journal_state",journal)
181 err=err+wfr_ji(w,"command_execution_reaped",reaped)
182 if reaped==1{err=err+wfr_ji(w,"command_exit_code",exitcode)}
183 err=err+wfr_js(w,"effects","connector-receipts-required")+wfr_js(w,"served_acceptance","unverified")
184 err=err+wfr_file_json(w,"evidence",evidence,log)+wfr_file_json(w,"ledger",ledger,f)
185 var truncated:i64=0;if log.status==1{if wf_has(log.text.data,"[NX-JOB CAPTURE-TRUNCATED")==1{truncated=1}}
186 err=err+wfr_ji(w,"evidence_truncation_marker",truncated)
187 err=err+json_emit_key(w,"runs",4)+json_begin_array(w)
188 var done:i64=0;var failed:i64=0;var parked:i64=0;var uncertain:i64=0;var i:i64=0
189 while i<nr{
190 let r:*WfrRun=((runs as i64)+i*__size_of(WfrRun)) as *WfrRun;var state:*u8="UNCERTAIN";var action:*u8="inspect-complete-evidence-and-reconcile-connector-receipts; do-not-auto-retry"
191 if frame==1{
192 if r.state==3{state="DONE";done=done+1;action="workflow-complete; review connector effects and served output"}
193 if r.state==2{state="FAILED";failed=failed+1}
194 if r.state==1{state="PARKED";parked=parked+1;action="inspect required decision; approve or deny explicitly"}
195 }
196 if seq(state,"UNCERTAIN")==1{uncertain=uncertain+1}
197 err=err+json_begin_object(w)+wfr_js(w,"run_id",r.rid)+wfr_run_labels(w,f.text.data,n,r.rid)+wfr_js(w,"state",state)
198 err=err+wfr_ji(w,"ok_steps",r.oks)+wfr_ji(w,"attempts",r.attempts)+wfr_ji(w,"failed_step",r.failedstep)+wfr_ji(w,"pending_step",r.pending)
199 err=err+wfr_js(w,"action",action)
200 err=err+json_emit_key(w,"delivery_report",15)+json_begin_object(w)
201 err=err+wfr_var_json(w,f.text.data,n,r.rid,"manifest")+wfr_var_json(w,f.text.data,n,r.rid,"manifest_sha256")+wfr_var_json(w,f.text.data,n,r.rid,"release_scope")+wfr_var_json(w,f.text.data,n,r.rid,"delivery_state")
202 err=err+json_end_object(w)+json_end_object(w);i=i+1
203 }
204 err=err+json_end_array(w)+wfr_ji(w,"selected_runs",nr)+wfr_ji(w,"done",done)+wfr_ji(w,"failed",failed)+wfr_ji(w,"parked",parked)+wfr_ji(w,"uncertain",uncertain)
205 var overall:*u8="needs-intervention"
206 if frame==1 && log.status==1 && truncated==0 && nr>0{if uncertain==0 && failed==0{if parked>0{overall="waiting-decision"}else{overall="workflow-complete-effects-unverified"}}}
207 if frame==1 && log.status==1 && truncated==0 && nr==0{overall="no-matching-runs";if mode==2{overall="no-eligible-runs"}}
208 if reaped==1 && exitcode!=0{overall="needs-intervention"}
209 err=err+wfr_js(w,"state",overall)+json_end_object(w)
210 if err<0{return 2}
211 if wfr_write_all(1,output.data,w.pos)!=0{return 2};if wfr_write_all(1,"
212",1)!=0{return 2}
213 if frame!=1 || log.status!=1 || truncated!=0{return 2};if reaped==1 && exitcode!=0{return 2};if failed>0 || uncertain>0{return 1};return 0
214}
215
216// Resume summarizes only runs that were eligible before this invocation.
217func wfr_resume_filter(ledger:*u8)->*WfOwnedText{
218 let f:*WfrFile=wfr_read(ledger)
219 if f.status!=1{return 0 as *WfOwnedText}
220 if wfr_frame(f.text.data,f.text.len)<0{return 0 as *WfOwnedText}
221 let out:*WfOwnedText=wf_text_new(f.text.len,1);if out==(0 as *WfOwnedText){return out}
222 var used:i64=0;var s:i64=0
223 while s<f.text.len{
224 var e:i64=s;while e<f.text.len{if f.text.data[e]==(10 as u8){break};e=e+1}
225 if wf_text_match(f.text.data,s,e,"WFRUN rid=",slen("WFRUN rid="))==1{
226 let st:*WfOwnedText=wfr_field(f.text.data,s,e," status=");if st==(0 as *WfOwnedText){return 0 as *WfOwnedText}
227 if seq(st.data,"START")==1{
228 let id:*WfOwnedText=wfr_field(f.text.data,s,e," rid=");if id==(0 as *WfOwnedText){return 0 as *WfOwnedText}
229 let r:*WfrRun=sys_mmap(__size_of(WfrRun)) as *WfrRun
230 if wfr_scan(f.text.data,f.text.len,id.data,r)!=0{return 0 as *WfOwnedText}
231 if r.state!=2 && r.state!=3{
232 if id.len>=out.bytes-used{return 0 as *WfOwnedText}
233 var j:i64=0;while j<id.len{out.data[used]=id.data[j];used=used+1;j=j+1};out.data[used]=10;used=used+1
234 };wf_text_free(id)
235 };wf_text_free(st)
236 };s=e+1
237 }
238 out.data[used]=0;out.len=used;return out
239}
240
241
242// Validate the complete event before effects. The engine's old first-ID/truncation
243// behavior is preserved for legacy verbs, but result verbs must not alias requests.
244func wfr_event_identity(ev:*u8)->*WfOwnedText{return wf_id_event(ev)}
245
246func wfr_refuse(reason:*u8,code:i64)->i64{
247 let keys:*u8="schema nishi-workflow-result/1 state refused effects none reason native_code action correct-input-before-dispatch"
248 let cap:i64=(slen(keys)+slen(reason))*6
249 let w:*JsonWriter=sys_mmap(__size_of(JsonWriter)) as *JsonWriter;let buf:*u8=sys_mmap(cap);json_writer_init(w,buf,cap)
250 var rc:i64=json_begin_object(w)
251 rc=rc+wfr_js(w,"schema","nishi-workflow-result/1")+wfr_js(w,"state","refused")+wfr_js(w,"effects","none")+wfr_js(w,"reason",reason)+wfr_ji(w,"native_code",code)+wfr_js(w,"action","correct-input-before-dispatch")+json_end_object(w)
252 if rc<0{return 2};wfr_write_all(1,buf,w.pos);wfr_write_all(1,"\n",1);return 2
253}
254
255func wfr_execute(argc:i64,argv:*i64,resume:i64)->i64{
256 var logindex:i64=7;if resume==1{logindex=6}
257 if argc<=logindex{return wfr_refuse("missing-evidence-path",0-1)}
258 var eventid:*WfOwnedText=0 as *WfOwnedText
259 if resume==0{eventid=wfr_event_identity(argv[6] as *u8);if eventid==(0 as *WfOwnedText){return wfr_refuse("invalid-or-aliased-event-identity",0-1)}}
260 var selection:*WfOwnedText=0 as *WfOwnedText
261 if resume==1{selection=wfr_resume_filter(argv[4] as *u8);if selection==(0 as *WfOwnedText){p("WFLOW result journal unavailable or malformed; no action started\n");return 2}}
262 let path:*u8=argv[logindex] as *u8
263 let fd:i64=sys_openat_exclusive(path,MODE_0600)
264 if fd<0{return wfr_refuse("evidence-path-exists-or-unavailable",fd)}
265 let shared:*i64=sys_mmap_shared(__size_of(i64)*2) as *i64
266 if (shared as i64)<=0{sys_close(fd);return 2};shared[0]=WF_EVIDENCE_ERROR;shared[1]=0
267 let pid:i64=sys_fork()
268 if pid<0{sys_close(fd);return 2}
269 if pid==0{
270 if fd!=1{if sys_dup3(fd,1,0)<0{sys_exit(2)}}
271 if fd!=2{if sys_dup3(fd,2,0)<0{sys_exit(2)}}
272 if fd>2{sys_close(fd)}
273 var rc:i64=0
274 if resume==1{rc=wf_resume_files(argv[2] as *u8,argv[3] as *u8,argv[4] as *u8,argv[5] as *u8)}else{rc=wf_fire_files(argv[2] as *u8,argv[3] as *u8,argv[4] as *u8,argv[5] as *u8,argv[6] as *u8)}
275 p("WFLOW-RESULT native_code=");pn(rc);p("
276")
277 shared[0]=rc
278 if sys_fsync(1)!=0{sys_exit(2)};shared[1]=1
279 if rc<0 && rc!=(0-100){sys_exit(1)};sys_exit(0)
280 }
281 sys_close(fd)
282 let status:*i64=sys_mmap(__size_of(i64)) as *i64
283 var waited:i64=sys_wait4(pid,status,0);while waited==(0-4){waited=sys_wait4(pid,status,0)}
284 var reaped:i64=0;var code:i64=0-1
285 if waited==pid{reaped=1;code=wait_status_rc(status[0])}
286 var filter:*u8="*";var mode:i64=0
287 if resume==1{filter=selection.data;mode=2}
288 if resume==0{filter=eventid.data;mode=1}
289 let result:i64=wfr_summary(argv[4] as *u8,path,filter,mode,reaped,code)
290 if reaped!=1 || code!=0 || shared[1]!=1{return 2};return result
291}