Skip to content

Commit 448e6f2

Browse files
authored
fix: Follow-up refactorings on ManagedResource class (#335)
* fix: Diagrams act as summary * fix: Shorten doc * fix: Shorten doc * fix: Shorten doc * fix: Shorten doc * fix: Shorten doc * fix: Shorten doc * fix: Shorten doc * fix: Shorten doc * fix: Add example in docs * fix: Add example in docs 2 * fix: Refactor * fix: Test fix
1 parent 6bed909 commit 448e6f2

2 files changed

Lines changed: 121 additions & 86 deletions

File tree

‎src/c2pa/c2pa.py‎

Lines changed: 95 additions & 84 deletions
Original file line numberDiff line numberDiff line change
@@ -255,6 +255,11 @@ class ManagedResource:
255255
The native pointer is freed automatically via `_free_native_ptr`.
256256
"""
257257

258+
_inflight = 0
259+
_mut_inflight = 0
260+
_pending_teardown: Optional[bool] = None
261+
_released = False
262+
258263
def _init_attrs(self):
259264
"""Set this class's own attributes to their defaults.
260265
@@ -333,7 +338,7 @@ def _ensure_not_borrowed(self):
333338
Raises:
334339
C2paError: If a native call is in flight on this resource.
335340
"""
336-
if getattr(self, '_inflight', 0) > 0:
341+
if self._inflight > 0:
337342
name = type(self).__name__
338343
raise C2paError(
339344
f"{name} is in use by another operation and "
@@ -345,7 +350,7 @@ def _ensure_no_mutating_call(self):
345350
Raises:
346351
C2paError: when a mutating native call is in progress.
347352
"""
348-
if getattr(self, '_mut_inflight', 0) > 0:
353+
if self._mut_inflight > 0:
349354
raise C2paError(
350355
f"{type(self).__name__} is running a mutating operation")
351356

@@ -371,45 +376,63 @@ def _guarded_op(self, *, refuse_mut=True):
371376
finally:
372377
self._maybe_flush_pending()
373378

374-
@contextlib.contextmanager
375-
def _native_call(self):
376-
"""Hold the handle valid across a native call that runs
377-
caller-supplied stream callbacks, so _live_op_lock() can't be held.
378-
Count the call as in-flight/in-progress.
379-
A free intent (teardown) is registered and the last caller frees.
380-
A free intent marks the resource as closed, preventing further use.
379+
def _begin_reservation(self, *, mutating, refuse):
380+
"""Count a native call as in flight, or raise.
381+
382+
Only place counters go up.
383+
`refuse` runs under the lock, before counting,
384+
and raises to reject the call.
381385
"""
382-
with self._live_op_lock():
383-
self._ensure_valid_state()
384-
self._ensure_no_mutating_call()
385-
self._inflight = getattr(self, '_inflight', 0) + 1
386386
try:
387-
with _native_section():
388-
yield
389-
finally:
390387
with self._live_op_lock():
391-
self._inflight -= 1
388+
self._ensure_valid_state()
389+
refuse()
390+
if mutating:
391+
self._mut_inflight += 1
392+
self._inflight += 1
393+
except BaseException:
392394
self._maybe_flush_pending()
395+
raise
396+
397+
def _end_reservation(self, *, mutating):
398+
"""Uncount a _begin_reservation() call and run any queued teardown.
399+
The only place the counters go down."""
400+
with self._live_op_lock():
401+
if mutating:
402+
self._mut_inflight -= 1
403+
self._inflight -= 1
404+
self._maybe_flush_pending()
393405

394406
@contextlib.contextmanager
395-
def _exclusive_native_call(self):
396-
"""Exclusively marks this handle as being mutated.
397-
A free intent (teardown) is registered and the last caller frees.
398-
A free intent marks the resource as closed, preventing further use.
407+
def _reserve(self, *, mutating, refuse):
408+
"""Hold a reservation across a native call.
409+
The lock is not held across the call itself:
410+
the call may run caller-supplied reentrant callbacks.
411+
A teardown that arrives meanwhile queues until
412+
the last call out.
399413
"""
400-
with self._live_op_lock():
401-
self._ensure_valid_state()
402-
self._ensure_no_mutating_call()
403-
self._mut_inflight = getattr(self, '_mut_inflight', 0) + 1
404-
self._inflight = getattr(self, '_inflight', 0) + 1
414+
self._begin_reservation(mutating=mutating, refuse=refuse)
405415
try:
406416
with _native_section():
407417
yield
408418
finally:
409-
with self._live_op_lock():
410-
self._mut_inflight -= 1
411-
self._inflight -= 1
412-
self._maybe_flush_pending()
419+
self._end_reservation(mutating=mutating)
420+
421+
@contextlib.contextmanager
422+
def _native_call(self):
423+
"""Reserve the handle for a shared call:
424+
other shared calls may run alongside it, no mutating call may."""
425+
with self._reserve(mutating=False,
426+
refuse=self._ensure_no_mutating_call):
427+
yield
428+
429+
@contextlib.contextmanager
430+
def _exclusive_native_call(self):
431+
"""Reserve the handle for a mutating call:
432+
no other mutating call may run alongside it."""
433+
with self._reserve(mutating=True,
434+
refuse=self._ensure_no_mutating_call):
435+
yield
413436

414437
@staticmethod
415438
def _free_native_ptr(ptr):
@@ -480,11 +503,10 @@ def _teardown(self, free_handle: bool):
480503
than report an error, so it cannot rely on acquiring.
481504
"""
482505
if is_foreign_process(self):
483-
self._handle = None
484-
self._lifecycle_state = LifecycleState.CLOSED
506+
self._detach_in_child()
485507
return
486508

487-
if getattr(self, '_released', False):
509+
if self._released:
488510
return
489511
self._record_pending_intent(free_handle)
490512

@@ -495,10 +517,10 @@ def _teardown(self, free_handle: bool):
495517
return
496518

497519
try:
498-
if getattr(self, '_released', False):
520+
if self._released:
499521
# Checks released as it recorded possible free intents.
500522
return
501-
if getattr(self, '_inflight', 0) > 0 or _in_native_section():
523+
if self._inflight > 0 or _in_native_section():
502524
# Closes the resource so it can't be used anymore.
503525
# Records also pending actual frees.
504526
self._close_lifecycle()
@@ -507,7 +529,7 @@ def _teardown(self, free_handle: bool):
507529
return
508530

509531
with self._live_teardown_lock():
510-
pending = getattr(self, '_pending_teardown', None)
532+
pending = self._pending_teardown
511533
if pending is not None:
512534
free_handle = pending and free_handle
513535
self._pending_teardown = None
@@ -523,7 +545,7 @@ def _record_pending_intent(self, free_handle: bool):
523545
free only leaks.
524546
"""
525547
with self._live_teardown_lock():
526-
pending = getattr(self, '_pending_teardown', None)
548+
pending = self._pending_teardown
527549
if pending is None:
528550
self._pending_teardown = free_handle
529551
else:
@@ -542,16 +564,25 @@ def _record_pending_teardown(self, free_handle: bool):
542564
self._record_pending_intent(free_handle)
543565
self._close_lifecycle()
544566

567+
def _detach_in_child(self):
568+
"""In a forked child:
569+
null this copy's handle and mark it closed, without freeing.
570+
The parent still owns the native pointer.
571+
"""
572+
if hasattr(self, '_handle'):
573+
self._handle = None
574+
if hasattr(self, '_lifecycle_state'):
575+
self._lifecycle_state = LifecycleState.CLOSED
576+
545577
def _finish_teardown(self, free_handle: bool):
546578
"""Once teardown can run, runs the actual release.
547579
Steps: release, null the handle, free if requested.
548580
"""
549581
if is_foreign_process(self):
550-
self._handle = None
551-
self._lifecycle_state = LifecycleState.CLOSED
582+
self._detach_in_child()
552583
return
553584

554-
if getattr(self, '_released', False):
585+
if self._released:
555586
# Already done by another caller (concurrent caller).
556587
return
557588

@@ -569,16 +600,16 @@ def _finish_teardown(self, free_handle: bool):
569600

570601
def _has_pending_teardown(self) -> bool:
571602
"""Check if a teardown request is waiting for the resource."""
572-
return getattr(self, '_pending_teardown', None) is not None
603+
return self._pending_teardown is not None
573604

574605
def _flush_pending_pass(self):
575606
"""Attempt to run pending teardowns.
576607
"""
577608

578609
with self._live_op_lock():
579-
if getattr(self, '_pending_teardown', None) is None:
610+
if self._pending_teardown is None:
580611
return
581-
if getattr(self, '_inflight', 0) > 0:
612+
if self._inflight > 0:
582613
return
583614
if _in_native_section():
584615
_register_for_section_flush(self)
@@ -596,8 +627,7 @@ def _maybe_flush_pending(self):
596627
return
597628

598629
self._flush_pending_pass()
599-
if self._has_pending_teardown() and not getattr(
600-
self, '_released', False):
630+
if self._has_pending_teardown() and not self._released:
601631
self._flush_pending_pass()
602632

603633
def _release_handle(self):
@@ -608,8 +638,7 @@ def _release_handle(self):
608638
to free.
609639
"""
610640
with self._live_op_lock():
611-
owned_elsewhere = getattr(
612-
self, '_pending_teardown', None) is not None
641+
owned_elsewhere = self._pending_teardown is not None
613642
if not owned_elsewhere and (
614643
self._lifecycle_state != LifecycleState.ACTIVE):
615644
self._handle = None
@@ -787,49 +816,36 @@ def _begin_consume(self):
787816
Raises:
788817
C2paError: Unusable resource or native call in progress.
789818
"""
790-
with self._live_op_lock():
791-
# A consumed or closed resource has no handle left to hand over;
792-
# without this the call would pass a null pointer to native.
793-
self._ensure_valid_state()
794-
self._ensure_not_borrowed()
795-
self._mut_inflight = getattr(self, '_mut_inflight', 0) + 1
796-
self._inflight = getattr(self, '_inflight', 0) + 1
819+
self._begin_reservation(mutating=True,
820+
refuse=self._ensure_not_borrowed)
797821

798822
def _end_consume(self):
799823
"""Release a _begin_consume() reservation, then run any teardown
800824
that arrived while it was held."""
801-
with self._live_op_lock():
802-
self._mut_inflight -= 1
803-
self._inflight -= 1
804-
self._maybe_flush_pending()
825+
self._end_reservation(mutating=True)
805826

806827
def _consume_and_swap(self, ffi_call, error_message):
807828
"""Run an FFI call consuming the handle, reserving it.
808829
A replacement handle will be swapping in on success
809830
(a returned null value is a failure).
810831
"""
832+
def swap(new_ptr):
833+
with self._live_op_lock():
834+
self._handle = new_ptr
811835

812-
self._begin_consume()
813-
try:
814-
with _native_section():
815-
new_ptr = self._invoke_consume(
816-
ffi_call, error_message, reserved=True)
817-
if new_ptr:
818-
with self._live_op_lock():
819-
self._handle = new_ptr
820-
return
821-
self._raise_consume_failure(error_message, reserved=True)
822-
finally:
823-
self._end_consume()
836+
self._consume_reserved(ffi_call, error_message,
837+
succeeded=bool, on_success=swap)
824838

825-
def _consume_reserved(self, ffi_call, error_message, *, succeeded):
826-
"""Run a reserved consuming call and mark the handle consumed on
827-
success.
839+
def _consume_reserved(self, ffi_call, error_message, *, succeeded,
840+
on_success=None):
841+
"""Run a reserved consuming call, then act on success.
828842
829843
Args:
830844
succeeded: Reads the call's raw result and returns whether it
831845
succeeded. Each entry point has its own convention: a status
832846
code, or a replacement pointer.
847+
on_success: Called with the raw result on success. Default:
848+
mark the handle consumed, closed, not freed.
833849
834850
Returns:
835851
The call's raw result, for callers that hand it on.
@@ -840,7 +856,10 @@ def _consume_reserved(self, ffi_call, error_message, *, succeeded):
840856
result = self._invoke_consume(
841857
ffi_call, error_message, reserved=True)
842858
if succeeded(result):
843-
self._teardown(free_handle=False)
859+
if on_success is None:
860+
self._teardown(free_handle=False)
861+
else:
862+
on_success(result)
844863
return result
845864
self._raise_consume_failure(error_message, reserved=True)
846865
finally:
@@ -896,15 +915,7 @@ def _cleanup_resources(self):
896915
"""Release native resources idempotently."""
897916
try:
898917
if is_foreign_process(self):
899-
# A forked child holds a separate copy of this object and the
900-
# parent still owns the real handle and frees it. Mark this
901-
# copy closed and null its handle so the child cannot mistake
902-
# it for usable or free it, but do not free here.
903-
# Mutating this copy does not touch the parent's.
904-
if hasattr(self, '_handle'):
905-
self._handle = None
906-
if hasattr(self, '_lifecycle_state'):
907-
self._lifecycle_state = LifecycleState.CLOSED
918+
self._detach_in_child()
908919
return
909920
if hasattr(self, '_lifecycle_state'):
910921
# Closes here must defer to the teardown checks.
@@ -921,7 +932,7 @@ def is_valid(self) -> bool:
921932
return (
922933
self._lifecycle_state == LifecycleState.ACTIVE
923934
and self._handle is not None
924-
and getattr(self, '_mut_inflight', 0) == 0
935+
and self._mut_inflight == 0
925936
)
926937

927938
def close(self) -> None:

‎tests/test_unit_tests_threaded.py‎

Lines changed: 26 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -4293,7 +4293,19 @@ def test_close_during_sign_does_not_deadlock(self):
42934293
"claim_generator_info": [
42944294
{"name": "python_test", "version": "0.0.1"}],
42954295
"format": "image/jpeg",
4296-
"assertions": [],
4296+
"assertions": [
4297+
{
4298+
"label": "c2pa.actions",
4299+
"data": {
4300+
"actions": [
4301+
{
4302+
"action": "c2pa.created",
4303+
"digitalSourceType": "http://cv.iptc.org/newscodes/digitalsourcetype/digitalCreation"
4304+
}
4305+
]
4306+
}
4307+
}
4308+
],
42974309
}
42984310
errors = []
42994311

@@ -5007,7 +5019,19 @@ def test_sign_with_internal_close_frees_once(self):
50075019
"claim_generator_info": [
50085020
{"name": "python_test", "version": "0.0.1"}],
50095021
"format": "image/jpeg",
5010-
"assertions": [],
5022+
"assertions": [
5023+
{
5024+
"label": "c2pa.actions",
5025+
"data": {
5026+
"actions": [
5027+
{
5028+
"action": "c2pa.created",
5029+
"digitalSourceType": "http://cv.iptc.org/newscodes/digitalsourcetype/digitalCreation"
5030+
}
5031+
]
5032+
}
5033+
}
5034+
],
50115035
}
50125036
signer = Signer.from_info(signer_info)
50135037
builder = Builder(manifest)

0 commit comments

Comments
 (0)