Aller au contenu

Référence Pipeline

Un pipeline enchaîne les étages qui transforment un texte en texte dé-identifié, et l'inverse. AnonymizationPipeline traite un seul texte, sans mémoire d'un appel à l'autre. ThreadAnonymizationPipeline traite une conversation, en gardant un jeton par valeur sur tous ses messages.

Les deux renvoient une Anonymization, le texte dé-identifié associé au jeton qui a remplacé chaque entité.


AnonymizationPipeline

Module : piighost.pipeline

Dé-identifie un seul texte à travers les étages. Dans l'ordre, ces étages détectent les données confidentielles, appliquent l'override du serveur, résolvent les spans qui se chevauchent, retrouvent les occurrences manquées, groupent les détections en entités, résolvent les conflits d'entités, remplacent par des jetons, puis revérifient avec un guard. Chaque appel à anonymize() est indépendant.

Constructeur

AnonymizationPipeline(
    detector: AnyDetector,
    linker: AnyEntityLinker | None = None,
    anonymizer: AnyAnonymizer[PreservationT] | None = None,
    overlap_resolver: AnyOverlapResolver | None = None,
    expander: AnyDetectionExpander | None = None,
    entity_resolver: AnyEntityResolver | None = None,
    guard: AnyGuardRail | None = None,
    observation_redactor: AnyPlaceholderFactory | None = None,
    override: AnyDetectionOverride | None = None,
    trace_clear_text: bool = False,
)
ParamètreTypeDéfautDescription
detectorAnyDetectorrequisDétecteur d'entités async
linkerAnyEntityLinker | NoneNoneGroupe les détections en entités. Par défaut ExactEntityLinker()
anonymizerAnyAnonymizer[PreservationT] | NoneNoneMoteur de remplacement et sa placeholder factory. Par défaut Anonymizer(LabelCounterPlaceholderFactory())
overlap_resolverAnyOverlapResolver | NoneNoneRésout les détections qui se chevauchent. Par défaut ConfidenceOverlapResolver(), car l'étape de rendu a besoin de spans disjoints
expanderAnyDetectionExpander | NoneNoneAjoute les occurrences manquées d'une valeur détectée. Désactivé quand None
entity_resolverAnyEntityResolver | NoneNoneRéconcilie les entités en conflit. Désactivé quand None
guardAnyGuardRail | NoneNoneRevérifie la sortie pour des valeurs confidentielles résiduelles. Désactivé quand None
observation_redactorAnyPlaceholderFactory | NoneNonePlaceholder factory remplaçant les valeurs en clair dans les payloads d'observation. Avec None, les traces gardent le texte en clair et peuvent ainsi servir de jeux d'annotation. Avec un tracer actif et sans masqueur, le constructeur émet un PIIGhostSecurityWarning sauf si trace_clear_text=True l'acquitte
overrideAnyDetectionOverride | NoneNoneListe à masquer et liste à laisser en clair du serveur imposées à chaque ensemble de détections. Désactivé quand None
trace_clear_textboolFalseAcquitte le traçage en clair de l'observation pour supprimer l'avertissement de sécurité quand aucun observation_redactor n'est défini

Méthodes

anonymize(text) -> Anonymization (async)

Exécute le pipeline complet et renvoie le texte dé-identifié avec le jeton utilisé pour chaque entité.

Lève PIIRemainingError quand un guard configuré signale des valeurs confidentielles restées dans la sortie.

result = await pipeline.anonymize("Patrick lives in Paris.")
# result.text == "<<PERSON:1>> lives in <<LOCATION:1>>."

deanonymize(text, tokens) -> str

Renvoie le texte avec chaque jeton connu remplacé par la valeur de son entité. tokens est la correspondance issue d'une Anonymization, lue à l'envers. Les jetons absents de la correspondance sont laissés intacts.

La restauration n'est sans ambiguïté que si les jetons préservent l'identité, car deux entités partageant un même jeton se confondent en une seule valeur.

original = pipeline.deanonymize(result.text, result.tokens)
# original == "Patrick lives in Paris."

ThreadAnonymizationPipeline

Module : piighost.pipeline

Dé-identifie chaque message d'une conversation avec des jetons stables sur toute la conversation. Une valeur vue dans un premier message puis à nouveau plus tard porte le même jeton, car les jetons sont assignés sur l'union des détections de tous les messages, pas sur un message seul. Les détections de chaque message sont mises en cache dans la mémoire. Un message renvoyé ne repasse donc pas par la détection.

Ce pipeline ajoute un composant, la mémoire de conversation memory. Elle stocke, pour chaque conversation, les détections de chaque message.

Constructeur

ThreadAnonymizationPipeline(
    detector: AnyDetector,
    linker: AnyEntityLinker | None = None,
    anonymizer: AnyAnonymizer[PreservationT] | None = None,
    memory: AnyConversationMemory | None = None,
    overlap_resolver: AnyOverlapResolver | None = None,
    expander: AnyDetectionExpander | None = None,
    entity_resolver: AnyEntityResolver | None = None,
    guard: AnyGuardRail | None = None,
    observation_redactor: AnyPlaceholderFactory | None = None,
    override: AnyDetectionOverride | None = None,
    trace_clear_text: bool = False,
    token_memo_ttl: float | None = None,
    time_source: Callable[[], float] = time.monotonic,
)

En plus de tous les paramètres de AnonymizationPipeline :

ParamètreTypeDéfautDescription
memoryAnyConversationMemory | NoneNoneStockage par conversation des détections de chaque message. Par défaut InMemoryConversationMemory() pour un seul processus. Passez RedisConversationMemory ou SqlAlchemyConversationMemory pour un backend partagé
token_memo_ttlfloat | NoneNoneSecondes pendant lesquelles la correspondance de jetons mémoïsée d'une conversation est gardée. Ce mémo garde les valeurs de la conversation en clair. forget_thread n'atteint que le processus où il tourne. Sur un déploiement multi-worker, ce délai borne donc combien de temps les autres workers gardent le mémo. None garde une entrée jusqu'à ce que la borne de taille l'évince
time_sourceCallable[[], float]time.monotonicL'horloge que lit token_memo_ttl, injectable pour les tests

Méthodes

anonymize(text, thread_id, role=MessageRole.USER) -> Anonymization (async)

Détecte les entités du message, les enregistre dans la mémoire de thread_id, puis dé-identifie avec des jetons assignés sur toute la conversation. Le jeton d'une valeur reste le même d'un message à l'autre.

Le thread_id est requis. Il n'y a pas de défaut partagé, donc deux appelants ne peuvent pas tomber dans une même conversation et laisser fuir mutuellement leurs données confidentielles. role marque l'auteur des valeurs que le message introduit. Une valeur introduite d'abord par l'assistant est laissée en clair, car ce n'est pas une donnée confidentielle de l'utilisateur.

Lève PIIRemainingError quand un guard configuré signale des valeurs confidentielles restées dans la sortie.

a1 = await pipeline.anonymize("Patrick lives in Paris.", thread_id="user-A")
a2 = await pipeline.anonymize("Patrick wrote to Marie.", thread_id="user-A")
# Patrick keeps <<PERSON:1>> across both turns.

anonymize_corrected(text, thread_id, detections) -> Anonymization (async)

Redé-identifie un message utilisateur avec un ensemble de détections corrigé par un humain. L'ensemble corrigé remplace les détections de ce message dans la mémoire, puis le message est dé-identifié avec des jetons cohérents sur la conversation. La détection ne relance pas. Cette méthode ne concerne que les propres messages d'un utilisateur. La correction est donc enregistrée comme un message utilisateur.

L'ensemble corrigé est stocké tel quel, sans résolution de chevauchement ni recherche d'occurrences, car l'humain fait autorité sur cet ensemble. Un override configuré s'applique encore, donc les listes du serveur priment sur la correction.

detection = Detection(span=Span(0, 5), text="Marie", label="PERSON", confidence=1.0)
detections = [detection]
result = await pipeline.anonymize_corrected("Marie called.", "user-A", detections)

deanonymize(text, thread_id) -> str (async)

Renvoie le texte avec chaque jeton de la conversation remplacé par sa valeur. Les jetons de la conversation sont reconstruits depuis sa mémoire. Tout texte qui les porte est donc restauré, y compris une réponse du modèle que le pipeline n'a jamais dé-identifiée.

reply = await pipeline.deanonymize("Message sent to <<PERSON:2>>.", thread_id="user-A")
# reply == "Message sent to Marie."

thread_token_map(thread_id) -> dict[str, str] (async)

Renvoie la correspondance placeholder vers valeur de la conversation, dérivée du cache. Avec elle, un appelant peut résoudre tout un flux d'un coup plutôt que de restaurer jeton par jeton. Un jeton que la conversation n'a jamais émis est absent de la correspondance.

forget_thread(thread_id) -> Forgotten (async)

Efface la mémoire d'une conversation et renvoie un Forgotten indiquant ce qui a été supprimé. Oublier une conversation inconnue ne supprime rien et rapporte zéro.

forgotten = await pipeline.forget_thread("user-A")
# forgotten.messages, forgotten.detections

Oublier une conversation efface aussi sa correspondance de jetons mémoïsée. Ce mémo garde les valeurs de la conversation en clair. Effacer le store seul laisserait donc ces valeurs vivantes dans le processus. Les autres conversations gardent leur mémo. L'appel n'atteint que le processus où il tourne. Sur un déploiement multi-worker, posez donc token_memo_ttl au constructeur pour borner la durée pendant laquelle les autres workers gardent le mémo, comme décrit dans Déploiement multi-instance. Un cache reste en place, celui des motifs de frontière de mot, partagé par tout le processus. Ce cache est indexé sur le fragment cherché, donc il contient des valeurs venant de toutes les conversations. Videz-le avec clear_boundary_cache quand une demande d'effacement couvre tout le processus.

from piighost.text import clear_boundary_cache

await pipeline.forget_thread("user-A")
clear_boundary_cache()

recognizer (propriété)

La grammaire des jetons que ce pipeline émet, une BaseDelimitedPlaceholderFactory, ou None. Une factory à délimiteurs est son propre recognizer, car ses jetons portent une grammaire retrouvable. Une factory sans grammaire, comme un masque, n'a pas de recognizer.


Ports

Deux protocoles typent un pipeline là où un appelant, comme le middleware, doit l'accepter sans dépendre d'une classe concrète. Les deux sont génériques sur ce que les jetons émis préservent. Un consommateur peut donc exiger un pipeline dont les jetons préservent l'identité, et rejeter un pipeline dont les jetons ne la préservent pas.

AnyPipeline

Un composant qui dé-identifie un seul texte et sait le restaurer.

@runtime_checkable
class AnyPipeline(Protocol[PreservationT_co]):
    async def anonymize(self, text: str) -> Anonymization[PreservationT_co]: ...
    def deanonymize(self, text: str, tokens: Mapping[Entity, str]) -> str: ...

AnyThreadPipeline

Un pipeline scopé par conversation, local ou distant. Il dé-identifie chaque message d'une conversation, redé-identifie un message corrigé, restaure tout texte portant les jetons de la conversation, oublie une conversation en entier, et expose la grammaire de ses jetons.

@runtime_checkable
class AnyThreadPipeline(Protocol[PreservationT_co]):
    async def anonymize(
        self, text: str, thread_id: str, role: MessageRole = MessageRole.USER
    ) -> Anonymization[PreservationT_co]: ...
    async def anonymize_corrected(
        self, text: str, thread_id: str, detections: list[Detection]
    ) -> Anonymization[PreservationT_co]: ...
    async def deanonymize(self, text: str, thread_id: str) -> str: ...
    async def forget_thread(self, thread_id: str) -> Forgotten: ...
    @property
    def recognizer(self) -> BaseDelimitedPlaceholderFactory | None: ...

BaseAnonymizationPipeline

Module : piighost.pipeline

La machinerie partagée que les deux pipelines étendent. Elle tient les composants d'étage et les étapes communes à tous les pipelines, c'est-à-dire les étages optionnels de chevauchement, de recherche d'occurrences et de résolution d'entités, la vérification du guard, et les payloads d'observation. Les pipelines concrets ajoutent leur propre anonymize, sur un texte seul ou sur une conversation.


Construire depuis une configuration

Module : piighost.config

load_pipeline et load_thread_pipeline lisent un fichier de configuration, TOML ou JSON selon son suffixe, ou une référence du catalogue, et renvoient un pipeline construit. Une configuration qui déclare une mémoire décrit un pipeline de conversation. Les deux loaders imposent cette distinction, et la vérifient avant de construire quoi que ce soit.

  • load_pipeline(path) renvoie un AnonymizationPipeline. Il lève ConfigError quand la configuration déclare une mémoire.
  • load_thread_pipeline(path) renvoie un ThreadAnonymizationPipeline. Il lève ConfigError quand la configuration ne déclare pas de mémoire.
from piighost.config import load_pipeline, load_thread_pipeline

pipeline = load_pipeline("pipeline.toml")
thread_pipeline = load_thread_pipeline("thread.toml")

Une référence écrite catalog:namespace/nom:sélecteur charge toute la configuration que le catalogue piighost publie sous ce nom, toutes les étapes comprises, exactement comme le ferait un fichier qui la contiendrait. Une référence épinglée sur un commit est téléchargée au premier chargement, puis lue depuis le cache disque. Une variable d'environnement préfixée PIIGHOST_ l'emporte sur une valeur du catalogue comme sur celle d'un fichier. Une référence écrite hub:namespace/nom:sélecteur, comme en 1.x, se charge de la même façon.

from piighost.config import load_pipeline

pipeline = load_pipeline("catalog:piighost/fr-notarial")

Ce package a besoin de l'extra config. Voir la référence Configuration TOML pour le format du fichier.


Exemple complet

import asyncio

from gliner2 import GLiNER2

from piighost.components.detector.ner.gliner2 import Gliner2Detector
from piighost.pipeline import ThreadAnonymizationPipeline

model = GLiNER2.from_pretrained("fastino/gliner2-multi-v1")
detector = Gliner2Detector(model=model, threshold=0.5, labels=["PERSON", "LOCATION"])
pipeline = ThreadAnonymizationPipeline(detector)


async def main():
    result = await pipeline.anonymize("Patrick is in Lyon.", thread_id="user-A")
    print(result.text)  # <<PERSON:1>> is in <<LOCATION:1>>.

    original = await pipeline.deanonymize(result.text, thread_id="user-A")
    print(original)  # Patrick is in Lyon.


asyncio.run(main())

Voir aussi