clusternode.java

来自「jsr170接口的java实现。是个apache的开源项目。」· Java 代码 · 共 1,278 行 · 第 1/3 页

JAVA
1,278
字号
     * @param register <code>true</code>, if this is a register operation;     *                 <code>false</code> otherwise     */    private void process(Collection c, boolean register) {        if (nodeTypeListener == null) {            String msg = "NodeType listener unavailable.";            log.error(msg);            return;        }        try {            if (register) {                nodeTypeListener.externalRegistered(c);            } else {                nodeTypeListener.externalUnregistered(c);            }        } catch (InvalidNodeTypeDefException e) {            String msg = "Unable to deliver node type operation: " + e.getMessage();            log.error(msg);        } catch (RepositoryException e) {            String msg = "Unable to deliver node type operation: " + e.getMessage();            log.error(msg);        }    }    /**     * Process a node type re-registration.     *     * @param ntDef node type definition     */    private void process(NodeTypeDef ntDef) {        if (nodeTypeListener == null) {            String msg = "NodeType listener unavailable.";            log.error(msg);            return;        }        try {            nodeTypeListener.externalReregistered(ntDef);        } catch (InvalidNodeTypeDefException e) {            String msg = "Unable to deliver node type operation: " + e.getMessage();            log.error(msg);        } catch (RepositoryException e) {            String msg = "Unable to deliver node type operation: " + e.getMessage();            log.error(msg);        }    }    /**     * Invoked when a record ends.     */    private void end() {        UpdateEventListener listener = null;        if (workspace != null) {            listener = (UpdateEventListener) wspUpdateListeners.get(workspace);            if (listener == null) {                try {                    clusterContext.updateEventsReady(workspace);                } catch (RepositoryException e) {                    String msg = "Error making update listener for workspace " +                            workspace + " online: " + e.getMessage();                    log.warn(msg);                }                listener = (UpdateEventListener) wspUpdateListeners.get(workspace);                if (listener ==  null) {                    String msg = "Update listener unavailable for workspace: " + workspace;                    log.error(msg);                    return;                }            }        } else {            if (versionUpdateListener != null) {                listener = versionUpdateListener;            } else {                String msg = "Version update listener unavailable.";                log.error(msg);                return;            }        }        try {            listener.externalUpdate(changeLog, events);        } catch (RepositoryException e) {            String msg = "Unable to deliver update events: " + e.getMessage();            log.error(msg);        }    }    //-------------------------------------------------------< RecordConsumer >    /**     * {@inheritDoc}     */    public String getId() {        return PRODUCER_ID;    }    /**     * {@inheritDoc}     */    public long getRevision() {        try {            return instanceRevision.get();        } catch (JournalException e) {            log.warn("Unable to return current revision.", e);            return Long.MAX_VALUE;        }    }    /**     * {@inheritDoc}     */    public void consume(Record record) {        log.info("Processing revision: " + record.getRevision());        String workspace = null;        try {            workspace = record.readString();            start(workspace);            for (;;) {                char c = record.readChar();                if (c == '\0') {                    break;                }                if (c == 'N') {                    NodeOperation operation = NodeOperation.create(record.readByte());                    operation.setId(record.readNodeId());                    process(operation);                } else if (c == 'P') {                    PropertyOperation operation = PropertyOperation.create(record.readByte());                    operation.setId(record.readPropertyId());                    process(operation);                } else if (c == 'E') {                    int type = record.readByte();                    NodeId parentId = record.readNodeId();                    Path parentPath = record.readPath();                    NodeId childId = record.readNodeId();                    Path.PathElement childRelPath = record.readPathElement();                    QName ntName = record.readQName();                    Set mixins = new HashSet();                    int mixinCount = record.readInt();                    for (int i = 0; i < mixinCount; i++) {                        mixins.add(record.readQName());                    }                    String userId = record.readString();                    process(createEventState(type, parentId, parentPath, childId,                            childRelPath, ntName, mixins, userId));                } else if (c == 'L') {                    NodeId nodeId = record.readNodeId();                    boolean isLock = record.readBoolean();                    if (isLock) {                        boolean isDeep = record.readBoolean();                        String owner = record.readString();                        process(nodeId, isDeep, owner);                    } else {                        process(nodeId);                    }                } else if (c == 'S') {                    String oldPrefix = record.readString();                    String newPrefix = record.readString();                    String uri = record.readString();                    process(oldPrefix, newPrefix, uri);                } else if (c == 'T') {                    int size = record.readInt();                    int opcode = size & NTREG_MASK;                    size &= ~NTREG_MASK;                    switch (opcode) {                        case NTREG_REGISTER:                            HashSet ntDefs = new HashSet();                            for (int i = 0; i < size; i++) {                                ntDefs.add(record.readNodeTypeDef());                            }                            process(ntDefs, true);                            break;                        case NTREG_REREGISTER:                            process(record.readNodeTypeDef());                            break;                        case NTREG_UNREGISTER:                            HashSet ntNames = new HashSet();                            for (int i = 0; i < size; i++) {                                ntNames.add(record.readQName());                            }                            process(ntNames, false);                            break;                        default:                            throw new IllegalArgumentException("Unknown opcode: " + opcode);                    }                } else {                    throw new IllegalArgumentException("Unknown entry type: " + c);                }            }            end();        } catch (JournalException e) {            String msg = "Unable to read revision '" + record.getRevision() + "'.";            log.error(msg, e);        } catch (IllegalArgumentException e) {            String msg = "Error while processing revision " +                    record.getRevision() + ": " + e.getMessage();            log.error(msg);        }    }    /**     * {@inheritDoc}     */    public void setRevision(long revision) {        try {            instanceRevision.set(revision);        } catch (JournalException e) {            log.warn("Unable to set current revision to " + revision + ".", e);        }    }    /**     * Create an event state.     *     * @param type event type     * @param parentId parent id     * @param parentPath parent path     * @param childId child id     * @param childRelPath child relative path     * @param ntName ndoe type name     * @param userId user id     * @return event     */    private EventState createEventState(int type, NodeId parentId, Path parentPath,                                        NodeId childId, Path.PathElement childRelPath,                                        QName ntName, Set mixins, String userId) {        switch (type) {            case Event.NODE_ADDED:                return EventState.childNodeAdded(parentId, parentPath, childId, childRelPath,                        ntName, mixins, getOrCreateSession(userId), true);            case Event.NODE_REMOVED:                return EventState.childNodeRemoved(parentId, parentPath, childId, childRelPath,                        ntName, mixins, getOrCreateSession(userId), true);            case Event.PROPERTY_ADDED:                return EventState.propertyAdded(parentId, parentPath, childRelPath,                        ntName, mixins, getOrCreateSession(userId), true);            case Event.PROPERTY_CHANGED:                return EventState.propertyChanged(parentId, parentPath, childRelPath,                        ntName, mixins, getOrCreateSession(userId), true);            case Event.PROPERTY_REMOVED:                return EventState.propertyRemoved(parentId, parentPath, childRelPath,                        ntName, mixins, getOrCreateSession(userId), true);            default:                String msg = "Unexpected event type: " + type;                throw new IllegalArgumentException(msg);        }    }    /**     * Return a session matching a certain user id.     *     * @param userId user id     * @return session     */    private Session getOrCreateSession(String userId) {        if (lastSession == null || !lastSession.getUserID().equals(userId)) {            lastSession = new ClusterSession(userId);        }        return lastSession;    }    //-----------------------------------------------< Record writing methods >    private static void write(Record record, ChangeLog changeLog, EventStateCollection esc)            throws JournalException {        Iterator deletedStates = changeLog.deletedStates();        while (deletedStates.hasNext()) {            ItemState state = (ItemState) deletedStates.next();            if (state.isNode()) {                write(record, NodeDeletedOperation.create((NodeState) state));            } else {                write(record, PropertyDeletedOperation.create((PropertyState) state));            }        }        Iterator modifiedStates = changeLog.modifiedStates();        while (modifiedStates.hasNext()) {            ItemState state = (ItemState) modifiedStates.next();            if (state.isNode()) {                write(record, NodeModifiedOperation.create((NodeState) state));            } else {                write(record, PropertyModifiedOperation.create((PropertyState) state));            }        }        Iterator addedStates = changeLog.addedStates();        while (addedStates.hasNext()) {            ItemState state = (ItemState) addedStates.next();            if (state.isNode()) {                write(record, NodeAddedOperation.create((NodeState) state));            } else {                write(record, PropertyAddedOperation.create((PropertyState) state));            }        }        Iterator events = esc.getEvents().iterator();        while (events.hasNext()) {            EventState event = (EventState) events.next();            write(record, event);        }    }    private static void write(Record record, String oldPrefix, String newPrefix, String uri)            throws JournalException {        record.writeChar('S');        record.writeString(oldPrefix);        record.writeString(newPrefix);        record.writeString(uri);    }    private static void write(Record record, Collection c, boolean register)            throws JournalException {        record.writeChar('T');        int size = c.size();        if (!register) {            size |= NTREG_UNREGISTER;        }        record.writeInt(size);        Iterator iter = c.iterator();        while (iter.hasNext()) {            if (register) {                record.writeNodeTypeDef((NodeTypeDef) iter.next());            } else {                record.writeQName((QName) iter.next());            }        }    }    private static void write(Record record, NodeTypeDef ntDef)            throws JournalException {        record.writeChar('T');        int size = 1;        size |= NTREG_REREGISTER;        record.writeInt(size);        record.writeNodeTypeDef(ntDef);    }    private static void write(Record record, PropertyOperation operation)            throws JournalException {        record.writeChar('P');        record.writeByte(operation.getOperationType());        record.writePropertyId(operation.getId());    }    private static void write(Record record, NodeOperation operation)            throws JournalException {        record.writeChar('N');        record.writeByte(operation.getOperationType());        record.writeNodeId(operation.getId());    }    /**     * Log an event. Subclass responsibility.     *     * @param event event to log     */    private static void write(Record record, EventState event)            throws JournalException {        record.writeChar('E');        record.writeByte(event.getType());        record.writeNodeId(event.getParentId());        record.writePath(event.getParentPath());        record.writeNodeId(event.getChildId());        record.writePathElement(event.getChildRelPath());        record.writeQName(event.getNodeType());        Set mixins = event.getMixinNames();        record.writeInt(mixins.size());        Iterator iter = mixins.iterator();        while (iter.hasNext()) {            record.writeQName((QName) iter.next());        }        record.writeString(event.getUserId());    }    /**     * Invoked when a cluster operation has ended. If <code>successful</code>,     * attempts to fill the journal record and update it, otherwise cancels     * the update.     *     * @param operation cluster operation     * @param successful <code>true</code> if the operation was successful and     *                   the journal record should be updated;     *                   <code>false</code> to revoke changes     */    public void ended(AbstractClusterOperation operation, boolean successful) {        Record record = operation.getRecord();        boolean succeeded = false;        try {            if (successful) {                record = operation.getRecord();                record.writeString(operation.getWorkspace());                operation.write();                record.writeChar('\0');                record.update();                setRevision(record.getRevision());                succeeded = true;            }        } catch (JournalException e) {            String msg = "Unable to create log entry: " + e.getMessage();            log.error(msg);        } catch (Throwable e) {            String msg = "Unexpected error while creating log entry.";            log.error(msg, e);        } finally {            if (!succeeded) {                record.cancelUpdate();            }        }    }}

⌨️ 快捷键说明

复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?